39#include "common/relpath.h"
40#include "utils/palloc.h"
67 char *rel = GetDatabasePath(db_oid, db_tablespace);
68 std::string path = std::string(DataDir) +
"/" + rel +
"/" + filename;
102 const std::vector<pg_uuid_t> &children)
104 const unsigned long idx =
gates.nbElements();
105 const unsigned long wires_idx =
wires.nbElements();
106 for(
const auto &c: children)
108 gates.add({type,
static_cast<unsigned>(children.size()), wires_idx});
126 &&
gates[idx].nb_children == 0;
127 bool real_create = type !=
gate_input || !children.empty();
128 if(placeholder && real_create) {
129 const unsigned long wires_idx =
wires.nbElements();
130 for(
const auto &c: children)
134 gates[idx].children_idx = wires_idx;
135 gates[idx].nb_children =
static_cast<unsigned>(children.size());
136 gates[idx].type = type;
137 for(
const auto &c: children)
146 for(
const auto &c: children)
157 return gates[idx].type;
162 std::vector<pg_uuid_t> result;
167 result.push_back(
wires[k]);
179 pg_uuid_t token,
double prob,
double *existing)
183 idx =
gates.nbElements();
192 if(std::isnan(prob)) {
193 gates[idx].prob = NAN;
198 gates[idx].prob = prob;
201 if(
gates[idx].prob == prob)
204 *existing =
gates[idx].prob;
215 double p =
gates[idx].prob;
218 if(
gates.version() < 2 && p == 1.)
228 if(!std::isnan(
gates[idx].prob))
229 return gates[idx].prob;
233 return 1. /
gates[idx].info2;
247 *prob =
gates[idx].prob;
259 pg_uuid_t token,
unsigned info1,
unsigned info2,
260 std::pair<unsigned, unsigned> *existing)
268 if((info1 != 0 && gi.
info1 != 0 && gi.
info1 != info1)
269 || (info2 != 0 && gi.
info2 != 0 && gi.
info2 != info2)) {
271 *existing = std::make_pair(gi.
info1, gi.
info2);
275 const unsigned w1 = info1 != 0 ? info1 : gi.
info1;
276 const unsigned w2 = info2 != 0 ? info2 : gi.
info2;
286 pg_uuid_t token,
const std::string &s, std::string *existing)
298 if(
gates[idx].extra_len == s.size()) {
300 for(
unsigned long k=0; k<s.size(); ++k)
309 if(
gates[idx].extra_len > 0) {
318 gates[idx].extra_len=s.size();
326 return std::make_pair(0, 0);
339 for(
unsigned long start=
gates[idx].extra_idx, k=start, end=start+
gates[idx].extra_len; k<end; ++k)
364 const char *names[4] = {
"provsql_mapping.mmap",
"provsql_gates.mmap",
365 "provsql_wires.mmap",
"provsql_extra.mmap" };
368 const bool committed = (access(marker.c_str(), F_OK) == 0);
370 for(
const char *name: names) {
373 if(access(fresh.c_str(), F_OK) != 0)
376 rename(fresh.c_str(), target.c_str());
378 unlink(fresh.c_str());
382 unlink(marker.c_str());
411 const std::vector<pg_uuid_t> &roots,
bool dry_run,
416 auto before = circuit->
counts();
421 std::vector<bool> live;
422 unsigned long nb_live = circuit->
mark(roots, live);
434 const char *names[4] = {
"provsql_mapping.mmap",
"provsql_gates.mmap",
435 "provsql_wires.mmap",
"provsql_extra.mmap" };
436 std::string target[4], fresh[4];
437 for(
int i=0; i<4; ++i) {
440 unlink(fresh[i].c_str());
443 auto after = circuit->
sweepInto(live, fresh[0], fresh[1], fresh[2], fresh[3]);
455 int mfd = open(marker.c_str(), O_CREAT | O_WRONLY | O_TRUNC, 0600);
457 const char *reason = strerror(errno);
459 marker.c_str(), reason);
464 for(
int i=0; i<4; ++i)
465 if(rename(fresh[i].c_str(), target[i].c_str())) {
466 const char *reason = strerror(errno);
468 target[i].c_str(), reason);
471 unlink(marker.c_str());
474#ifdef PROVSQL_INPROCESS_STORE
486 const std::vector<pg_uuid_t> &roots)
503 if(c==
'C' || c==
'G' || c==
'P' || c==
'I' || c==
'E')
511 unsigned nb_children;
516 std::vector<pg_uuid_t> children(nb_children);
517 for(
unsigned i=0; i<nb_children; ++i)
534 unsigned nb_children, info1, info2, len;
540 std::vector<pg_uuid_t> children(nb_children);
541 for(
unsigned i=0; i<nb_children; ++i)
544 if(!
READM(has_infos,
char) || !
READM(info1,
unsigned) || !
READM(info2,
unsigned)
545 || !
READM(len,
unsigned))
547 std::vector<char> data(len);
553 std::pair<unsigned, unsigned> existing{0, 0};
554 if(circuit->
setInfos(token, info1, info2, &existing)
556 provsql_warning(
"gate %s already records the annotation (%u, %u), not (%u, %u)",
557 uuid2string(token).c_str(), existing.first, existing.second,
561 std::string existing;
562 if(circuit->
setExtra(token, std::string(data.data(), len), &existing)
564 provsql_warning(
"gate %s already records the annotation \"%s\", not \"%s\"",
566 std::string(data.data(), len).c_str());
574 double prob, existing = 0.;
579 auto result = circuit->
setProb(token, prob, &existing);
580 char return_value =
static_cast<char>(result);
582 if(!
WRITEB(&return_value,
char) || !
WRITEB(&existing,
double))
583 provsql_error(
"Cannot write response to pipe (message type P)");
598 char has = circuit->
hasProb(token, &prob) ? 1 : 0;
601 provsql_error(
"Cannot write response to pipe (message type q)");
618 unsigned info1, info2;
619 std::pair<unsigned, unsigned> existing{0, 0};
620 if(!
READM(info1,
unsigned) || !
READM(info2,
unsigned))
622 if(circuit->
setInfos(token, info1, info2, &existing)
624 provsql_warning(
"gate %s already records the annotation (%u, %u), not (%u, %u)",
625 uuid2string(token).c_str(), existing.first, existing.second,
629 if(!
READM(len,
unsigned))
632 std::vector<char> data(len);
633 std::string existing;
636 if(circuit->
setExtra(token, std::string(data.data(), len), &existing)
638 provsql_warning(
"gate %s already records the annotation \"%s\", not \"%s\"",
640 std::string(data.data(), len).c_str());
656 provsql_error(
"Cannot write response to pipe (message type t)");
664 if(!
WRITEB(&nb,
unsigned long))
665 provsql_error(
"Cannot write response to pipe (message type n)");
677 unsigned nb_children = children.size();
678 if(!
WRITEB(&nb_children,
unsigned))
679 provsql_error(
"Cannot write response to pipe (message type c)");
682 provsql_error(
"Cannot write response to pipe (message type c)");
693 double prob = circuit->
getProb(token);
695 if(!
WRITEB(&prob,
double))
696 provsql_error(
"Cannot write response to pipe (message type p)");
707 auto infos = circuit->
getInfos(token);
709 if(!
WRITEB(&infos.first,
unsigned) || !
WRITEB(&infos.second,
unsigned))
710 provsql_error(
"Cannot write response to pipe (message type i)");
721 auto str = circuit->
getExtra(token);
722 unsigned len = str.size();
725 provsql_error(
"Cannot write response to pipe (message type e)");
736#ifdef PROVSQL_INPROCESS_STORE
739 provsql_error(
"message type g is not used by the in-process store");
741 std::stringstream ss;
742 boost::archive::binary_oarchive oa(ss);
745 ss.seekg(0, std::ios::end);
746 unsigned long size = ss.tellg();
747 ss.seekg(0, std::ios::beg);
768 unsigned long nb_roots;
770 if(!
READM(dry_run,
char) || !
READM(nb_roots,
unsigned long))
773 std::vector<pg_uuid_t> roots(nb_roots);
774 for(
unsigned long i=0; i<nb_roots; ++i)
779 cleanupStore(db_oid, db_tablespace, roots, dry_run != 0, &res);
784 provsql_error(
"Cannot write response to pipe (message type X)");
798 provsql_error(
"Cannot write response to pipe (message type S)");
806 char unclean = chk.
unclean ? 1 : 0;
807 if(!
WRITEB(&unclean,
char)
815 provsql_error(
"Cannot write response to pipe (message type k)");
827 if(!
READM(nb_roots,
unsigned))
830 std::vector<pg_uuid_t> roots(nb_roots);
831 for(
unsigned i=0; i<nb_roots; ++i)
835#ifdef PROVSQL_INPROCESS_STORE
838 provsql_error(
"message type j is not used by the in-process store");
840 std::stringstream ss;
841 boost::archive::binary_oarchive oa(ss);
844 ss.seekg(0, std::ios::end);
845 unsigned long size = ss.tellg();
846 ss.seekg(0, std::ios::beg);
859#ifndef PROVSQL_INPROCESS_STORE
876 int r = poll(&pfd, 1,
891 Oid db_oid, db_tablespace;
892 if(!
READM(db_oid, Oid) || !
READM(db_tablespace, Oid))
924 return mapping.uncleanShutdown() ||
gates.uncleanShutdown()
925 ||
wires.uncleanShutdown() ||
extra.uncleanShutdown();
936 std::vector<bool> referenced(c.
nb_gates,
false);
937 for(
unsigned long k=0; k<
mapping.capacity(); ++k) {
938 unsigned long v =
mapping.slotValue(k);
944 referenced[v] =
true;
946 for(
unsigned long i=0; i<c.
nb_gates; ++i) {
960 std::vector<bool> &live)
const
962 const unsigned long n =
gates.nbElements();
963 live.assign(n,
false);
964 unsigned long nb_live = 0;
966 std::vector<unsigned long> stack;
967 for(
const auto &r: roots) {
972 stack.push_back(idx);
976 while(!stack.empty()) {
977 unsigned long idx = stack.back();
987 stack.push_back(child);
996 const std::vector<bool> &live,
997 const std::string &mp,
const std::string &gp,
998 const std::string &wp,
const std::string &ep)
const
1001 const unsigned long n =
gates.nbElements();
1007 unsigned long next = 0;
1008 for(
unsigned long i=0; i<n; ++i)
1022 for(
unsigned long i=0; i<n; ++i) {
1037 nwires.
add(
wires[old_children + c]);
1039 const unsigned long old_extra = gi.
extra_idx;
1054 if(legacy && gi.
prob == 1.
1062 for(
unsigned long k=0; k<
mapping.capacity(); ++k) {
1063 unsigned long v =
mapping.slotValue(k);
1090 return memcmp(&a, &b,
sizeof(
pg_uuid_t))<0;
1099 const std::vector<pg_uuid_t> &roots)
const
1107 std::set<pg_uuid_t> to_process, processed;
1108 for(
const auto &r : roots)
1109 to_process.insert(r);
1113 while(!to_process.empty()) {
1115 to_process.erase(to_process.begin());
1116 processed.insert(uuid);
1122 if(!std::isnan(prob))
1125 std::vector<pg_uuid_t> children =
getChildren(uuid);
1126 for(
unsigned i=0; i<children.size(); ++i) {
1130 if(processed.find(children[i])==processed.end())
1131 to_process.insert(children[i]);
1136 auto [info1, info2] =
getInfos(uuid);
1145 auto [info1, info2] =
getInfos(uuid);
1146 if(info1 != 0 || info2 != 0)
gate_t
Strongly-typed gate identifier.
Out-of-line template method implementations for Circuit<gateType>.
Semiring-agnostic in-memory provenance circuit.
static constexpr const char * CLEANUP_SUFFIX
Suffix of the files a clean-up builds beside the live ones.
static void finishInterruptedCleanup(Oid db_oid, Oid db_tablespace)
Bring a database's store to a definite state before opening it.
static void cleanupStore(Oid db_oid, Oid db_tablespace, const std::vector< pg_uuid_t > &roots, bool dry_run, provsql_cleanup_result *out)
Rebuild a database's store, keeping only what roots reach.
static bool carriesProb(gate_type type)
Whether type is one of the gate kinds that carry a probability.
static std::map< Oid, MMappedCircuit * > circuits
Per-database mmap-backed provenance circuits, keyed by database OID.
void destroy_provsql_mmap()
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.
bool operator<(const pg_uuid_t a, const pg_uuid_t b)
Lexicographic less-than comparison for pg_uuid_t.
void provsql_mmap_main_loop()
Main processing loop of the mmap background worker.
static MMappedCircuit * getCircuit(Oid db_oid, Oid db_tablespace)
Return (creating lazily if needed) the circuit for db_oid.
void initialize_provsql_mmap()
Initialise the circuit store.
static bool store_dirty
Whether any circuit has been written to since the last flush.
static constexpr const char * CLEANUP_MARKER
Marker file present only while a clean-up's rename sequence is in flight; see finishInterruptedCleanu...
Persistent, mmap-backed storage for the full provenance circuit.
void addWire(gate_t f, gate_t t)
Add a directed wire from gate f (parent) to gate t (child).
gate_t getGate(const uuid &u)
Return (or create) the gate associated with UUID u.
In-memory provenance circuit with semiring-generic evaluation.
void setInfos(gate_t g, unsigned info1, unsigned info2)
Set the integer annotation pair for gate g.
gate_t setGate(gate_type type) override
Allocate a new gate with type type and no UUID.
void setExtra(gate_t g, const std::string &ex)
Attach a string extra to gate g.
void setProb(gate_t g, double p)
Set the probability for gate g.
Persistent mmap-backed representation of the provenance circuit.
static constexpr uint16_t GATES_VERSION
Format version of the gates file this build writes.
unsigned long mark(const std::vector< pg_uuid_t > &roots, std::vector< bool > &live) const
Mark every gate reachable from roots.
SetAnnotationResult setInfos(pg_uuid_t token, unsigned info1, unsigned info2, std::pair< unsigned, unsigned > *existing=nullptr)
Write the info1 / info2 annotations of a gate, once.
MMappedUUIDHashTable mapping
UUID → gate-index hash table.
void createGate(pg_uuid_t token, gate_type type, const std::vector< pg_uuid_t > &children)
Persist a new gate to the mmap store.
Counts counts() const
This store's current size.
std::string getExtra(pg_uuid_t token) const
Return the variable-length string annotation for gate token.
bool legacyProbabilities() const
Whether this store's gates file predates the NaN unset-probability convention (see GATES_VERSION).
SetProbResult setProb(pg_uuid_t token, double prob, double *existing=nullptr)
Write a gate's probability, once.
void appendGate(pg_uuid_t token, gate_type type, const std::vector< pg_uuid_t > &children)
Append a complete gate record, then publish token for it.
unsigned long getNbGates() const
Return the total number of gates stored in the circuit.
SetProbResult
Outcome of setProb.
@ NotProbGate
The token names a gate that carries no probability.
@ Unchanged
The gate already held exactly this probability.
@ Written
The gate had no probability and now holds the given one.
@ AlreadySet
The gate holds a different probability; refused.
static constexpr const char * GATES_FILENAME
Backing file for gates.
Counts sweepInto(const std::vector< bool > &live, const std::string &mp, const std::string &gp, const std::string &wp, const std::string &ep) const
Copy the gates flagged in live into a fresh set of files beside the current ones (suffix "....
gate_type getGateType(pg_uuid_t token) const
Return the type of the gate identified by token.
Check check() const
Walk the store and report what does not add up.
void sync()
Flush all backing files to disk with msync().
static std::string storePath(Oid db_oid, Oid db_tablespace, const char *filename)
Build the full path of a file in a database's store directory, the one GetDatabasePath resolves for d...
static constexpr const char * WIRES_FILENAME
Backing file for wires.
MMappedCircuit(const std::string &mp, const std::string &gp, const std::string &wp, const std::string &ep, bool read_only)
Delegating constructor that accepts pre-built paths.
GenericCircuit createGenericCircuit(pg_uuid_t token) const
Build an in-memory GenericCircuit rooted at token.
bool uncleanShutdown() const
Whether any backing file was found still marked open-for-writing when it was opened – the previous wr...
static constexpr const char * EXTRA_FILENAME
Backing file for extra.
static constexpr uint64_t MAGIC_WIRES
static constexpr uint64_t MAGIC_GATES
8-byte magic constants identifying each mmap file type.
static constexpr const char * MAPPING_FILENAME
Backing file for mapping.
bool hasProb(pg_uuid_t token, double *prob) const
Report whether token names a gate that holds a probability.
MMappedVector< char > extra
Variable-length string data.
static constexpr uint64_t MAGIC_EXTRA
double getProb(pg_uuid_t token) const
Return the probability an evaluation would use for token.
std::vector< pg_uuid_t > getChildren(pg_uuid_t token) const
Return the child UUIDs of the gate identified by token.
MMappedVector< GateInformation > gates
Gate metadata array.
SetAnnotationResult setExtra(pg_uuid_t token, const std::string &s, std::string *existing=nullptr)
Attach a variable-length string annotation to a gate, once.
SetAnnotationResult
Outcome of setInfos or setExtra.
@ Unchanged
The gate already held exactly this annotation.
@ Written
The gate had none and now holds the given annotation.
@ NoSuchGate
The token names no gate.
@ AlreadySet
The gate holds a different annotation; refused.
bool hasProbAt(unsigned long idx) const
Whether the gate record at index idx holds a written probability (see GATES_VERSION for the version-1...
MMappedVector< pg_uuid_t > wires
Flattened child UUID array.
static constexpr uint64_t MAGIC_MAPPING
std::pair< unsigned, unsigned > getInfos(pg_uuid_t token) const
Return the info1 / info2 pair for the gate token.
void flush()
Force every backing file to stable storage.
Persistent open-addressing hash table mapping UUIDs to integers.
std::pair< unsigned long, bool > publish(pg_uuid_t u, unsigned long value)
Insert UUID u with a caller-chosen value.
static constexpr unsigned long NOTHING
Sentinel returned by operator[]() when the UUID is not present.
void flush()
Force the backing file to stable storage (MappedRegion::flush()).
Append-only, mmap-backed vector of elements of type T.
void flush()
Force the backing file to stable storage (MappedRegion::flush()).
unsigned long nbElements() const
Return the number of elements currently stored.
void add(const T &value)
Append an element to the end of the vector.
#define provsql_error(fmt,...)
Report a fatal ProvSQL error and abort the current transaction.
#define provsql_warning(fmt,...)
Emit a ProvSQL warning message (execution continues).
bool provsql_worker_buffered(void)
Whether the worker's read buffer holds unread bytes, which poll() cannot see.
Background worker and IPC primitives for mmap-backed circuit storage.
#define WRITEB_BYTES(ptr, n)
Write n reply bytes to the main-to-background pipe.
#define READM(var, type)
Read one value of type from the background-to-main pipe.
#define WRITEB(pvar, type)
Write one value of type to the main-to-background pipe.
#define PROVSQL_STORE_FLUSH_INTERVAL_MS
How long the worker waits after a write before forcing the store to stable storage.
#define READM_BYTES(ptr, n)
Read exactly n bytes of a request from the background-to-main pipe.
provsqlSharedState * provsql_shared_state
Pointer to the ProvSQL shared-memory segment (set in provsql_shmem_startup).
Shared-memory segment and inter-process pipe management.
@ gate_observe
Latent-variable observation (likelihood-weighting evidence): one wire → an observed bare gate_rv leaf...
@ gate_rv
Continuous random-variable leaf (extra encodes distribution).
@ gate_annotation
Transparent single-child wrapper carrying a query-level annotation in extra (inversion-free certifica...
@ gate_mobius
Signed Möbius combination: a MEASURE-only gate carrying one integer coefficient per child (in extra,...
@ gate_arith
n-ary arithmetic gate over scalar-valued children (info1 holds operator tag)
@ gate_assumed
Structural marker over a single child whose sub-circuit was computed under a Boolean-provenance assum...
string uuid2string(pg_uuid_t uuid)
Format a pg_uuid_t as a std::string.
C++ utility functions for UUID manipulation.
What provsql.check_store() reports about a store.
unsigned long dangling_indices
Mapping entries indexing past the records.
unsigned long nb_gates
Gate records.
unsigned long bad_extra
Records whose extra runs past the extra file.
bool unclean
A file was left marked open-for-writing.
unsigned long nb_mapping
Mapping entries.
unsigned long unreferenced
Records no mapping entry points at.
unsigned long bad_wires
Records whose children run past the wires.
unsigned long next_value
Next index the mapping would assign.
The size of a store, in the three units that matter.
unsigned long extra_bytes
Bytes of variable-length annotation.
unsigned long gates
Gate records.
unsigned long wires
Child wires.
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.