ProvSQL C/C++ API
Adding support for provenance and uncertainty management to PostgreSQL databases
Loading...
Searching...
No Matches
circuit_cleanup.c
Go to the documentation of this file.
1/**
2 * @file circuit_cleanup.c
3 * @brief Removing from the circuit store what nothing references any more.
4 *
5 * The store only grows. A gate is never removed, and that is deliberate:
6 * gates are immutable and content-addressed, so a transaction that rolls
7 * back leaves an orphan rather than an inconsistency, and the same
8 * expression recomputed lands on the same gate. The price is that every
9 * dropped table, every @c remove_provenance, every reloaded dataset and
10 * every exploratory join over a large table leaves gates behind for good.
11 *
12 * @c provsql.circuit_cleanup() is the complement of that: the one
13 * operation allowed to remove gates, run explicitly, the way
14 * @c VACUUM @c FULL is the complement of MVCC. It keeps every gate
15 * reachable from a token stored in the database and discards the rest,
16 * rewriting the four store files compactly on the way -- which also makes
17 * it the repair tool for a store damaged by an interrupted write (see
18 * @c provsql.check_store).
19 *
20 * **Why it needs the database to itself.** Two properties rule out a
21 * concurrent collector. Gates are re-created idempotently: a query
22 * running alongside can compute @c provenance_plus(a, b), find the
23 * content-addressed gate already there as an orphan, and reference it a
24 * moment before the collector -- which saw no reference -- removes it.
25 * And a transaction in progress has created gates and inserted rows that
26 * the collector's snapshot cannot see. So the function takes the
27 * database's lock in the mode @c DROP @c DATABASE takes it, and refuses
28 * to run while another session is connected, exactly as @c DROP
29 * @c DATABASE does.
30 *
31 * **What counts as a root.** Every value of a @c uuid, @c agg_token or
32 * @c random_variable column -- and of arrays of those -- in every table
33 * and materialised view of the database, whatever the column is called:
34 * @c "CREATE TABLE ... AS SELECT provenance()" stores tokens under
35 * whatever name the user picks, mapping tables use @c provenance,
36 * @c update_provenance uses @c provsql, and users store @c get_children
37 * results and conditioning events wherever they like. Scanning by type
38 * is the only conservative choice, because a missed root turns a derived
39 * gate back into an input, silently. The semiring constants @c gate_zero
40 * and @c gate_one are roots unconditionally: the planner emits their
41 * UUIDs as literals, so nothing in the database need mention them.
42 *
43 * A token that lives only *outside* the database -- copied into a
44 * notebook cell or a deep link, kept in a file, stored as text or inside
45 * a @c jsonb document -- is not a root. Content-addressed gates come
46 * back by re-running the query that produced them; freshly minted ones
47 * (an @c rv leaf, the @c update gate of a deleted log row, the input gate
48 * of a row deleted from an untracked copy) do not.
49 */
50#include "postgres.h"
51
52#include "access/htup_details.h"
53#include "catalog/pg_database.h"
54#include "commands/dbcommands.h"
55#include "executor/spi.h"
56#include "fmgr.h"
57#include "funcapi.h"
58#include "miscadmin.h"
59#include "storage/lmgr.h"
60#include "storage/procarray.h"
61#include "utils/builtins.h"
62#include "utils/guc.h"
63#include "utils/lsyscache.h"
64#include "utils/uuid.h"
65
66#include "circuit_cache.h"
67#include "provsql_mmap.h"
68#include "provsql_shmem.h"
69#include "provsql_utils.h"
70
71/** @brief Growable array of root tokens gathered from the database. */
72typedef struct root_set {
74 int64 len;
75 int64 cap;
76} root_set;
77
78static void root_set_add(root_set *rs, const pg_uuid_t *token)
79{
80 if(rs->len == rs->cap) {
81 int64 newcap = rs->cap ? rs->cap * 2 : 1024;
82 pg_uuid_t *grown = repalloc(rs->tokens, newcap * sizeof(pg_uuid_t));
83 rs->tokens = grown;
84 rs->cap = newcap;
85 }
86 rs->tokens[rs->len++] = *token;
87}
88
89/**
90 * @brief Run @p query and append every UUID it returns to @p rs.
91 *
92 * The query is expected to return one @c uuid column.
93 */
94static void collect_from(root_set *rs, const char *query)
95{
96 int rc = SPI_execute(query, true, 0);
97 if(rc != SPI_OK_SELECT)
98 provsql_error("circuit_cleanup: cannot collect roots (SPI code %d)", rc);
99
100 for(uint64 i = 0; i < SPI_processed; ++i) {
101 bool isnull;
102 Datum d = SPI_getbinval(SPI_tuptable->vals[i], SPI_tuptable->tupdesc,
103 1, &isnull);
104 if(!isnull)
105 root_set_add(rs, DatumGetUUIDP(d));
106 }
107}
108
109/**
110 * @brief Gather every token stored in the database.
111 *
112 * Walks @c pg_attribute for columns whose type is one ProvSQL tokens live
113 * in, then reads each of them. Provenance tracking is switched off for
114 * the duration: these reads would otherwise materialise an input gate per
115 * row of every tracked table, which is exactly the growth the clean-up is
116 * there to undo.
117 */
118static void collect_roots(root_set *rs)
119{
120 int save_nestlevel = NewGUCNestLevel();
121 int rc;
122
123 SetConfigOption("provsql.active", "off",
124 PGC_USERSET, PGC_S_SESSION);
125
126 PG_TRY();
127 {
128 /* The constants (zero, one, the NULL value): the planner emits their
129 UUIDs as literals, so no row need mention them, and a store that lost
130 them would read them back as inputs rather than as constants. */
131 {
132 const char *constants[] = {PROVSQL_GATE_ZERO_UUID, PROVSQL_GATE_ONE_UUID,
134 for(int i = 0; i < 3; ++i)
135 root_set_add(rs, DatumGetUUIDP(
136 DirectFunctionCall1(uuid_in,
137 CStringGetDatum(constants[i]))));
138 }
139
140 rc = SPI_execute(
141 "SELECT c.oid::regclass::text AS rel, quote_ident(a.attname) AS col, "
142 " t.typelem <> 0 AND t.typlen = -1 AS is_array "
143 " FROM pg_attribute a "
144 " JOIN pg_class c ON c.oid = a.attrelid "
145 " JOIN pg_namespace n ON n.oid = c.relnamespace "
146 " JOIN pg_type t ON t.oid = a.atttypid "
147 " WHERE c.relkind IN ('r', 'm', 'p') "
148 " AND a.attnum > 0 AND NOT a.attisdropped "
149 " AND n.nspname NOT IN ('pg_catalog', 'information_schema', 'pg_toast') "
150 " AND (CASE WHEN t.typelem <> 0 AND t.typlen = -1 THEN t.typelem "
151 " ELSE t.oid END) IN ("
152 " 'uuid'::regtype, 'provsql.agg_token'::regtype, "
153 " 'provsql.random_variable'::regtype) "
154 " ORDER BY 1, 2", true, 0);
155 if(rc != SPI_OK_SELECT)
156 provsql_error("circuit_cleanup: cannot enumerate token columns "
157 "(SPI code %d)", rc);
158
159 {
160 uint64 n = SPI_processed;
161 char **rels = palloc(n * sizeof(char *));
162 char **cols = palloc(n * sizeof(char *));
163 bool *arrs = palloc(n * sizeof(bool));
164
165 for(uint64 i = 0; i < n; ++i) {
166 rels[i] = SPI_getvalue(SPI_tuptable->vals[i], SPI_tuptable->tupdesc, 1);
167 cols[i] = SPI_getvalue(SPI_tuptable->vals[i], SPI_tuptable->tupdesc, 2);
168 arrs[i] = (strcmp(SPI_getvalue(SPI_tuptable->vals[i],
169 SPI_tuptable->tupdesc, 3), "t") == 0);
170 }
171
172 for(uint64 i = 0; i < n; ++i) {
173 StringInfoData buf;
174 initStringInfo(&buf);
175 if(arrs[i])
176 appendStringInfo(&buf,
177 "SELECT DISTINCT u::uuid FROM %s, "
178 "LATERAL unnest(%s) AS u WHERE u IS NOT NULL",
179 rels[i], cols[i]);
180 else
181 appendStringInfo(&buf,
182 "SELECT DISTINCT %s::uuid FROM %s WHERE %s IS NOT NULL",
183 cols[i], rels[i], cols[i]);
184 collect_from(rs, buf.data);
185 pfree(buf.data);
186 }
187 }
188
189 /* A foreign table can hold tokens too, but scanning one means talking
190 to another server -- which may be down, slow, or not there at all.
191 Say so rather than silently treating its rows as absent. */
192 rc = SPI_execute(
193 "SELECT count(*) FROM pg_attribute a "
194 " JOIN pg_class c ON c.oid = a.attrelid "
195 " JOIN pg_type t ON t.oid = a.atttypid "
196 " WHERE c.relkind = 'f' AND a.attnum > 0 AND NOT a.attisdropped "
197 " AND (CASE WHEN t.typelem <> 0 AND t.typlen = -1 THEN t.typelem "
198 " ELSE t.oid END) IN ("
199 " 'uuid'::regtype, 'provsql.agg_token'::regtype, "
200 " 'provsql.random_variable'::regtype)", true, 0);
201 if(rc == SPI_OK_SELECT && SPI_processed == 1) {
202 char *cnt = SPI_getvalue(SPI_tuptable->vals[0], SPI_tuptable->tupdesc, 1);
203 if(cnt && strcmp(cnt, "0") != 0)
204 provsql_notice("circuit_cleanup: %s token-typed column(s) live on "
205 "foreign tables and were not scanned; gates only they "
206 "reference are removed", cnt);
207 }
208 }
209 PG_CATCH();
210 {
211 AtEOXact_GUC(false, save_nestlevel);
212 PG_RE_THROW();
213 }
214 PG_END_TRY();
215
216 AtEOXact_GUC(false, save_nestlevel);
217}
218
219PG_FUNCTION_INFO_V1(circuit_cleanup);
220/**
221 * @brief Rebuild this database's circuit store, keeping only what the
222 * tokens stored in the database reach.
223 *
224 * Takes the database exclusively for the duration; see the file comment
225 * for why that is unavoidable. @p dry_run reports how much would be kept
226 * without writing anything.
227 */
228Datum circuit_cleanup(PG_FUNCTION_ARGS)
229{
230 bool dry_run = PG_ARGISNULL(0) ? false : PG_GETARG_BOOL(0);
231 root_set rs = { NULL, 0, 0 };
233 int nprepared = 0;
234 TupleDesc tupdesc;
235 Datum values[6];
236 bool nulls[6] = {false, false, false, false, false, false};
237
238 if(RecoveryInProgress())
239 ereport(ERROR,
240 (errmsg("provsql.circuit_cleanup() cannot run on a standby"),
241 errdetail("The circuit store is not replicated, and a standby "
242 "must not write to it.")));
243
245 ereport(ERROR,
246 (errmsg("provsql.circuit_cleanup() cannot run in a transaction "
247 "that has already written to the circuit store"),
248 errhint("Run it as the first statement of its own transaction.")));
249
250 /* The lock DROP DATABASE takes: every new connection takes it in a
251 weaker mode during its startup transaction, so sessions that arrive
252 from now on wait for us. */
253 LockSharedObject(DatabaseRelationId, MyDatabaseId, 0, AccessExclusiveLock);
254
255 /* ... and the check DROP DATABASE makes, for the sessions that are
256 already here. CountOtherDBBackends terminates autovacuum workers and
257 waits for them; anything else is the operator's to disconnect. */
258 {
259 int nbackends = 0;
260 if(CountOtherDBBackends(MyDatabaseId, &nbackends, &nprepared))
261 ereport(ERROR,
262 (errcode(ERRCODE_OBJECT_IN_USE),
263 errmsg("database \"%s\" is being accessed by other users",
264 get_database_name(MyDatabaseId)),
265 errdetail("There %s %d other session%s and %d prepared "
266 "transaction%s using the database.",
267 nbackends == 1 ? "is" : "are", nbackends,
268 nbackends == 1 ? "" : "s", nprepared,
269 nprepared == 1 ? "" : "s")));
270 }
271
272 /* Our own caches answer "this gate exists" without contacting the
273 worker, which would survive the rebuild as a lie. */
276
277 /* The root set outlives the SPI session that fills it, so it is
278 allocated in our own context: SPI_finish frees everything palloc'd
279 while connected. */
280 rs.tokens = MemoryContextAlloc(CurrentMemoryContext,
281 1024 * sizeof(pg_uuid_t));
282 rs.cap = 1024;
283
284 if(SPI_connect() != SPI_OK_CONNECT)
285 provsql_error("circuit_cleanup: cannot connect to SPI");
286 collect_roots(&rs);
287 SPI_finish();
288
289 provsql_circuit_cleanup_request(dry_run, rs.tokens, rs.len, &res);
290
291 if(get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE)
292 provsql_error("circuit_cleanup: expected composite return type");
293 tupdesc = BlessTupleDesc(tupdesc);
294
295 values[0] = Int64GetDatum((int64) res.gates_before);
296 values[1] = Int64GetDatum((int64) res.gates_after);
297 values[2] = Int64GetDatum((int64) res.wires_before);
298 values[3] = Int64GetDatum((int64) res.wires_after);
299 values[4] = Int64GetDatum((int64) res.extra_before);
300 values[5] = Int64GetDatum((int64) res.extra_after);
301 if(dry_run) {
302 /* A dry run measures the mark phase only: the wire and byte totals of
303 a rewrite are not known without doing it. */
304 nulls[3] = true;
305 nulls[5] = true;
306 }
307
308 PG_RETURN_DATUM(HeapTupleGetDatum(heap_form_tuple(tupdesc, values, nulls)));
309}
C-linkage interface to the in-process provenance circuit cache.
void circuit_cache_reset(void)
Forget every cached gate.
static void collect_roots(root_set *rs)
Gather every token stored in the database.
Datum circuit_cleanup(PG_FUNCTION_ARGS)
Rebuild this database's circuit store, keeping only what the tokens stored in the database reach.
static void root_set_add(root_set *rs, const pg_uuid_t *token)
static void collect_from(root_set *rs, const char *query)
Run query and append every UUID it returns to rs.
void provsql_gate_builders_forget(void)
Forget what gate_builders.c remembers of the store (planted gates, value gates written).
#define provsql_error(fmt,...)
Report a fatal ProvSQL error and abort the current transaction.
#define provsql_notice(fmt,...)
Emit a ProvSQL informational notice (execution continues).
bool provsql_store_written(void)
Whether this transaction has written to the circuit store.
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.
Background worker and IPC primitives for mmap-backed circuit storage.
Shared-memory segment and inter-process pipe management.
Core types, constants, and utilities shared across ProvSQL.
#define PROVSQL_GATE_ZERO_UUID
UUID of the semiring zero gate: the result of the SQL function gate_zero(), created with the extensio...
#define PROVSQL_GATE_ONE_UUID
UUID of the semiring one gate, the result of gate_one().
#define PROVSQL_GATE_NULL_UUID
UUID of the constant value gate standing for the NULL value, the result of gate_null().
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.
Growable array of root tokens gathered from the database.
pg_uuid_t * tokens