ProvSQL C/C++ API
Adding support for provenance and uncertainty management to PostgreSQL databases
Loading...
Searching...
No Matches
provsql_mmap.h
Go to the documentation of this file.
1/**
2 * @file provsql_mmap.h
3 * @brief Background worker and IPC primitives for mmap-backed circuit storage.
4 *
5 * ProvSQL persists the provenance circuit in memory-mapped files so that
6 * data survives transaction boundaries and is shared across backend
7 * processes. Because multiple backends may create gates concurrently, a
8 * dedicated PostgreSQL background worker (@c provsql_mmap_worker) is the
9 * sole writer to those files; normal backends communicate with it through
10 * a pair of anonymous pipes described in @c provsqlSharedState.
11 *
12 * This header exposes:
13 * - Functions to register, start, and manage the background worker.
14 * - A set of pipe I/O macros (@c READM, @c READB, @c WRITEB, @c WRITEM)
15 * that wrap @c read()/@c write() calls on the inter-process pipes.
16 * - A buffered-write interface (@c STARTWRITEM, @c ADDWRITEM, @c SENDWRITEM)
17 * that batches multiple fields into a single @c write() to stay within
18 * the atomic @c PIPE_BUF guarantee.
19 */
20#ifndef PROVSQL_MMAP_H
21#define PROVSQL_MMAP_H
22
23#include "limits.h"
24#include <unistd.h>
25
26#include "postgres.h"
27#include "provsql_utils.h"
28#include "provsql_config.h"
29
30/**
31 * @brief Entry point for the ProvSQL mmap background worker.
32 *
33 * Called by the postmaster when it launches the background worker.
34 * Enters the main loop (@c provsql_mmap_main_loop()) and never returns
35 * normally. The single @c Datum argument is required by the
36 * background-worker API but is not used.
37 */
38void provsql_mmap_worker(Datum);
39
40/**
41 * @brief Register the ProvSQL mmap background worker with PostgreSQL.
42 *
43 * Must be called from the extension's @c _PG_init() function so that
44 * the postmaster starts the worker on the next connection.
45 */
47
48/**
49 * @brief Initialise the circuit store.
50 *
51 * Called once by the background worker at startup. The per-database
52 * files themselves are opened lazily, on the first message for their
53 * database.
54 */
56
57/**
58 * @brief Unmap and close the mmap files.
59 *
60 * Called by the background worker on shutdown to release resources and
61 * ensure all dirty pages are synced to disk via @c msync().
62 */
63void destroy_provsql_mmap(void);
64
65/**
66 * @brief Main processing loop of the mmap background worker.
67 *
68 * Waits for gate-creation requests from backend processes, processes them
69 * by writing to the mmap files, and handles SIGTERM for graceful shutdown.
70 */
71void provsql_mmap_main_loop(void);
72
73/**
74 * @brief How long the worker waits after a write before forcing the store
75 * to stable storage.
76 *
77 * The circuit store is outside PostgreSQL's WAL, so a committed
78 * transaction's gates can still be sitting in the kernel's page cache
79 * when the machine loses power. Forcing them out this long after the
80 * last write bounds the loss, the way @c synchronous_commit @c = @c off
81 * bounds the heap's; @c provsql.synchronous_commit removes it entirely,
82 * at the price of one flush per store-writing transaction.
83 */
84#define PROVSQL_STORE_FLUSH_INTERVAL_MS 200
85
86/**
87 * @brief What @c provsql.circuit_cleanup() reports.
88 *
89 * The three "before" figures are the store's size when the clean-up
90 * started; the three "after" figures are the size of the rebuilt store
91 * (or, on a dry run, only the live gate count, the rest being 0 -- the
92 * wire and byte totals of a rewrite are not known without doing it).
93 */
94typedef struct provsql_cleanup_result {
95 uint64 gates_before; ///< Gate records before
96 uint64 gates_after; ///< Gate records kept
97 uint64 wires_before; ///< Child wires before
98 uint64 wires_after; ///< Child wires kept
99 uint64 extra_before; ///< Annotation bytes before
100 uint64 extra_after; ///< Annotation bytes kept
102
103/**
104 * @brief Ask the worker to rebuild this database's store, keeping only
105 * what @p roots reach.
106 *
107 * The caller must hold the database exclusively; see
108 * @c provsql.circuit_cleanup in @c circuit_cleanup.c.
109 */
110void provsql_circuit_cleanup_request(bool dry_run, const pg_uuid_t *roots,
111 int64 nb_roots,
113
114/** @brief Force every open circuit to stable storage. */
115void provsql_store_flush(void);
116
117/**
118 * @brief Note that this transaction has written to the circuit store.
119 *
120 * Arms the at-commit sync barrier (@c provsql.synchronous_commit) and the
121 * @c PREPARE @c TRANSACTION refusal.
122 */
123void provsql_store_note_write(void);
124
125/** @brief Whether this transaction has written to the circuit store. */
126bool provsql_store_written(void);
127
128/**
129 * @brief Handle a single IPC message: read its payload and write its reply.
130 *
131 * The opcode @p c and the message header (@p db_oid and the database's
132 * default tablespace @p db_tablespace) have already been consumed by the
133 * caller. Shared by the background-worker main loop (multi-process build)
134 * and the synchronous in-process dispatcher.
135 */
136void provsql_mmap_dispatch(char c, Oid db_oid, Oid db_tablespace);
137
138/**
139 * @brief Create a gate from in-extension C/C++ code (cache + worker IPC).
140 *
141 * Internal entry point behind the SQL-callable @c create_gate(), without
142 * Datum marshalling or gate-type-OID lookups; idempotent on
143 * already-mapped tokens.
144 *
145 * @param token UUID of the gate.
146 * @param type Gate type.
147 * @param nb_children Number of children.
148 * @param children_data Child UUIDs (may be NULL when @p nb_children is 0).
149 */
150void provsql_internal_create_gate(const pg_uuid_t *token, gate_type type,
151 unsigned nb_children,
152 const pg_uuid_t *children_data);
153
154/**
155 * @brief Create a gate together with its infos and its text, in one message
156 * that is not answered.
157 *
158 * For the gates whose address determines what they record (a value gate and
159 * its text, an annotation, a comparison and its operator, an aggregate), so
160 * that the write-once rule has nothing to refuse and nobody needs to wait for
161 * its answer. @p extra NULL records no text; @p has_infos false records no
162 * infos.
163 */
165 unsigned nb_children,
166 const pg_uuid_t *children,
167 bool has_infos, unsigned info1,
168 unsigned info2, const char *extra);
169
170/**
171 * @brief Outcome of a probability write, mirroring
172 * @c MMappedCircuit::SetProbResult across the IPC boundary.
173 */
175 PROVSQL_SET_PROB_NOT_PROB_GATE = 0, ///< The gate carries no probability
176 PROVSQL_SET_PROB_WRITTEN = 1, ///< Written; undo on rollback
177 PROVSQL_SET_PROB_UNCHANGED = 2, ///< Already held exactly this value
178 PROVSQL_SET_PROB_ALREADY_SET = 3 ///< Holds a different value; refused
180
181/**
182 * @brief Write a gate's probability from in-extension C/C++ code.
183 *
184 * Probabilities are written once (see @c MMappedCircuit::setProb), so
185 * this reports which of the four cases applied rather than a bare
186 * success flag. Callers that write a probability on a gate they have
187 * just created can treat anything but @c PROVSQL_SET_PROB_NOT_PROB_GATE
188 * as success; @c set_prob() itself raises on
189 * @c PROVSQL_SET_PROB_ALREADY_SET.
190 *
191 * Note that this is the raw store operation: it does not record the
192 * write in the transaction's undo list. SQL-level writers go through
193 * @c provsql_set_prob_tracked() in @c probability_store.c so a rollback
194 * drops what they wrote.
195 *
196 * @param token UUID of the gate.
197 * @param prob Probability value in [0,1], or @c NaN to clear.
198 * @param existing On @c PROVSQL_SET_PROB_ALREADY_SET, the stored value.
199 */
201 double prob,
202 double *existing);
203
204/**
205 * @brief Drop a gate's probability, leaving it as it was before anyone
206 * wrote one.
207 *
208 * The rollback path of @c probability_store.c, and the only way to unset
209 * a probability: there is none from SQL. Unlike a write it does not arm
210 * the at-commit sync barrier -- it runs when the transaction that would
211 * have committed is already gone.
212 */
213void provsql_internal_clear_prob(const pg_uuid_t *token);
214
215/**
216 * @brief Report whether a probability has been written on a gate.
217 *
218 * Distinct from @c get_prob(), which reports the value an evaluation
219 * would use and so answers 1 for a gate nobody gave a probability.
220 *
221 * @param token UUID of the gate.
222 * @param prob On @c true return, the written probability.
223 */
224bool provsql_internal_get_prob_written(const pg_uuid_t *token, double *prob);
225
226/**
227 * @brief Fetch a gate's type and children, cache-first with a worker
228 * round-trip (and cache fill) on a miss.
229 *
230 * On return @p *children_out is a @c calloc'd array to be freed by the
231 * caller, or @c NULL when the gate has no children.
232 */
234 unsigned *nb_children_out,
235 pg_uuid_t **children_out);
236
237
238#ifdef PROVSQL_INPROCESS_STORE
239
240/**
241 * @brief In-process replacement for a pipe write of a complete request.
242 *
243 * Appends the message in @p buf (@p len bytes) to the request FIFO and runs
244 * @c provsql_mmap_dispatch once, leaving any reply in the response FIFO for
245 * the caller's @c READB / @c READB_BYTES to consume.
246 */
247bool provsql_inproc_send(const char *buf, size_t len);
248
249/** Growable shared write buffer used with @c STARTWRITEM / @c ADDWRITEM. */
250extern char *buffer;
251/** Current write position within @c buffer. */
252extern unsigned bufferpos;
253/** Allocated capacity of @c buffer. */
254extern size_t buffercap;
255/** @brief Ensure @c buffer can hold at least @p need bytes. */
256void provsql_buffer_ensure(size_t need);
257
258#define READM(var, type) provsql_fifo_pop (&provsql_shared_state->req, &(var), sizeof(type))
259#define READB(var, type) provsql_fifo_pop (&provsql_shared_state->resp, &(var), sizeof(type))
260#define WRITEB(pvar, type) provsql_fifo_push(&provsql_shared_state->resp, (pvar), sizeof(type))
261#define WRITEM(pvar, type) provsql_fifo_push(&provsql_shared_state->req, (pvar), sizeof(type))
262
263#define READB_BYTES(ptr, n) provsql_fifo_pop (&provsql_shared_state->resp, (ptr), (n))
264#define READM_BYTES(ptr, n) provsql_fifo_pop (&provsql_shared_state->req, (ptr), (n))
265#define WRITEB_BYTES(ptr, n) provsql_fifo_push(&provsql_shared_state->resp, (ptr), (n))
266
267#define STARTWRITEM() (bufferpos=0)
268#define ADDWRITEM(pvar, type) (provsql_buffer_ensure(bufferpos+sizeof(type)), memcpy(buffer+bufferpos, pvar, sizeof(type)), bufferpos+=sizeof(type))
269#define SENDWRITEM() provsql_inproc_send(buffer, bufferpos)
270
271#else
272
273/** @brief Read exactly @p n bytes from @p fd into @p dst; @c false on EOF/error. */
274bool provsql_read_all(int fd, void *dst, size_t n);
275
276/** Shared write buffer used with @c STARTWRITEM / @c ADDWRITEM / @c SENDWRITEM. */
277extern char buffer[PIPE_BUF];
278/** Current write position within @c buffer. */
279extern unsigned bufferpos;
280
281/**
282 * @brief The worker's buffered read of the request pipe: what the pipe
283 * holds is read in one call, and the messages are parsed from the
284 * buffer. @c false on EOF or error.
285 */
286bool provsql_worker_read(void *dst, size_t n);
287/** @brief Whether the worker's read buffer holds unread bytes, which
288 * @c poll() cannot see. */
289bool provsql_worker_buffered(void);
290
291/** @brief Read one value of @p type from the background-to-main pipe. */
292#define READM(var, type) provsql_worker_read(&var, sizeof(type))
293/** @brief Read one value of @p type from the main-to-background pipe. */
294#define READB(var, type) (read(provsql_shared_state->pipembr, &var, sizeof(type))==(ssize_t)sizeof(type)) // flawfinder: ignore
295/** @brief Write one value of @p type to the main-to-background pipe. */
296#define WRITEB(pvar, type) (write(provsql_shared_state->pipembw, pvar, sizeof(type))!=-1)
297/** @brief Write one value of @p type to the background-to-main pipe. */
298#define WRITEM(pvar, type) (write(provsql_shared_state->pipebmw, pvar, sizeof(type))!=-1)
299
300/** @brief Read exactly @p n bytes of a reply from the main-to-background pipe. */
301#define READB_BYTES(ptr, n) provsql_read_all(provsql_shared_state->pipembr, (ptr), (n)) // flawfinder: ignore
302/** @brief Read exactly @p n bytes of a request from the background-to-main pipe. */
303#define READM_BYTES(ptr, n) provsql_worker_read((ptr), (n))
304/** @brief Write @p n reply bytes to the main-to-background pipe. */
305#define WRITEB_BYTES(ptr, n) (write(provsql_shared_state->pipembw, (ptr), (n))!=-1)
306
307/** @brief Reset the shared write buffer for a new batched write. */
308#define STARTWRITEM() (bufferpos=0)
309/** @brief Append one value of @p type to the shared write buffer. */
310#define ADDWRITEM(pvar, type) (memcpy(buffer+bufferpos, pvar, sizeof(type)), bufferpos+=sizeof(type))
311/** @brief Flush the shared write buffer to the background-to-main pipe atomically. */
312#define SENDWRITEM() (write(provsql_shared_state->pipebmw, buffer, bufferpos)!=-1)
313
314#endif /* PROVSQL_INPROCESS_STORE */
315
316/**
317 * @brief Append the per-message database header to the write buffer.
318 *
319 * Every request carries the OID of the database it applies to and the
320 * OID of that database's default tablespace, so the worker -- which runs
321 * outside any transaction and cannot read @c pg_database -- can resolve
322 * the directory holding the backing files. Follows the opcode byte in
323 * every message.
324 */
325#define ADDWRITEDB() (ADDWRITEM(&MyDatabaseId, Oid), ADDWRITEM(&MyDatabaseTableSpace, Oid))
326
327#endif /* PROVSQL_COLUMN_NAME */
Build-configuration switches shared across the C and C++ sources.
char buffer[PIPE_BUF]
Shared write buffer used with STARTWRITEM / ADDWRITEM / SENDWRITEM.
unsigned bufferpos
Current write position within buffer.
void initialize_provsql_mmap(void)
Initialise the circuit store.
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)
Create a gate from in-extension C/C++ code (cache + worker IPC).
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.
provsql_set_prob_result
Outcome of a probability write, mirroring MMappedCircuit::SetProbResult across the IPC boundary.
@ PROVSQL_SET_PROB_NOT_PROB_GATE
The gate carries no probability.
@ PROVSQL_SET_PROB_WRITTEN
Written; undo on rollback.
@ PROVSQL_SET_PROB_ALREADY_SET
Holds a different value; refused.
@ PROVSQL_SET_PROB_UNCHANGED
Already held exactly this value.
void provsql_mmap_worker(Datum)
Entry point for the ProvSQL mmap background worker.
bool provsql_read_all(int fd, void *dst, size_t n)
Read exactly n bytes from fd into dst; false on EOF/error.
void destroy_provsql_mmap(void)
Unmap and close the mmap files.
void provsql_store_flush(void)
Force every open circuit to stable storage.
void provsql_mmap_dispatch(char c, Oid db_oid, Oid db_tablespace)
Handle a single IPC message: read its payload and write its reply.
gate_type provsql_fetch_gate(const pg_uuid_t *token, unsigned *nb_children_out, pg_uuid_t **children_out)
Fetch a gate's type and children, cache-first with a worker round-trip (and cache fill) on a miss.
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,...
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.
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.
void provsql_mmap_main_loop(void)
Main processing loop of the mmap background worker.
void RegisterProvSQLMMapWorker(void)
Register the ProvSQL mmap background worker with PostgreSQL.
Core types, constants, and utilities shared across ProvSQL.
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.