Everett
Loading...
Searching...
No Matches
world.h
Go to the documentation of this file.
1
13#pragma once
14
15#include "everett/pins.h"
16
17#include <algorithm>
18#include <array>
19#include <atomic>
20#include <concepts>
21#include <cstdint>
22#include <istream>
23#include <limits>
24#include <map>
25#include <memory>
26#include <optional>
27#include <ostream>
28#include <span>
29#include <stdexcept>
30#include <string>
31#include <string_view>
32#include <utility>
33#include <vector>
34
35namespace everett {
36 template <class V> struct world_record {
37 std::string key;
38 std::optional<V> value;
39 bool operator==(world_record const &) const = default;
40 };
41
42 template <class V> struct world_edit {
43 std::string key;
44 std::optional<V> before;
45 std::optional<V> after;
46 bool operator==(world_edit const &) const = default;
47 };
48
49 template <class V> struct world_run {
50 explicit world_run(std::vector<world_record<V>> records)
51 : records_(std::move(records)), object_id_(next_identity()) {}
52
53 std::span<world_record<V> const> records() const noexcept { return records_; }
54 std::string const & object_id() const noexcept { return object_id_; }
55
56 private:
57 static std::string next_identity() {
58 // Reference runs have process-local identities. Durable blob stores must
59 // assign stable object identities of their own; signatures never do so.
60 static std::atomic<std::uint64_t> sequence = 0;
61 return "reference-run-" + std::to_string(sequence.fetch_add(1, std::memory_order_relaxed));
62 }
63
64 std::vector<world_record<V>> records_;
65 std::string object_id_;
66 };
67
68 namespace world_detail {
69 inline void write_u64(std::ostream & out, std::uint64_t value) {
70 std::array<char, 8> bytes;
71 for (unsigned i = 0; i < 8; ++i) bytes[i] = char((value >> (8 * i)) & 255);
72 out.write(bytes.data(), bytes.size());
73 if (!out) throw std::runtime_error("world debug export write failed");
74 }
75
76 inline std::uint64_t read_u64(std::istream & in) {
77 std::array<unsigned char, 8> bytes;
78 in.read(reinterpret_cast<char *>(bytes.data()), bytes.size());
79 if (!in) throw std::runtime_error("truncated world debug export");
80 std::uint64_t value = 0;
81 for (unsigned i = 0; i < 8; ++i) value |= std::uint64_t(bytes[i]) << (8 * i);
82 return value;
83 }
84 }
85
87 // The codec tag identifies values in debug resolved-table dumps.
88 // Custom fixed-width values supply a codec.
89 static constexpr std::uint64_t format_tag = 1;
90 static void write(std::ostream & out, std::uint64_t value) {
91 world_detail::write_u64(out, value);
92 }
93 static std::uint64_t read(std::istream & in) { return world_detail::read_u64(in); }
94 };
95
96 // Allocation bounds for import_rc of a resolved-table dump.
98 std::uint64_t max_records = 1'000'000;
99 std::uint64_t max_key_bytes = 64 * 1024 * 1024;
100 };
101
102 template <class V, class A, class H, class P> struct partition_round;
103
104 // Semantic reference implementation. A snapshot pins immutable sorted runs,
105 // and updates copy only the pin vector, never earlier record arrays. This has
106 // no redundant COLA scheduling or front-coded disk layout yet. It is an oracle
107 // for that implementation, not a claim of its asymptotic space/time bounds.
108 template <class V = std::uint64_t, class A = wrapping_fingerprint_algebra,
109 class H = u64_table_hash>
111 using value_type = V;
112 using element = typename A::element;
115 using pin = std::shared_ptr<run const>;
117
118 explicit reference_world(H hash = {})
119 : state_(std::make_shared<state>()), hash_(std::move(hash)) {}
120
121 static reference_world from_records(std::vector<record> records, H hash = {}) {
122 std::sort(records.begin(), records.end(),
123 [](record const & a, record const & b) { return a.key < b.key; });
124 auto result = reference_world(std::move(hash));
125 auto next = std::make_shared<state>();
126 auto signature = A::zero();
127 for (std::size_t i = 0; i < records.size(); ++i) {
128 if (!records[i].value) throw std::invalid_argument("initial record is a tombstone");
129 if (i && records[i - 1].key == records[i].key)
130 throw std::invalid_argument("duplicate initial world key");
131 signature = A::add(signature,
132 table_fingerprint<A>::binding(records[i].key, records[i].value, result.hash_));
133 }
134 next->live_size = records.size();
135 if (!records.empty()) {
136 auto object = std::make_shared<run>(std::move(records));
137 next->owner = next->owner.add({object->object_id(), object, signature, signature});
138 }
139 result.state_ = std::move(next);
140 return result;
141 }
142
143 std::optional<V> get(std::string_view key) const {
144 auto objects = state_->owner.pins();
145 for (auto it = objects.rbegin(); it != objects.rend(); ++it) {
146 auto entries = (*it)->records();
147 auto found = std::lower_bound(entries.begin(), entries.end(), key,
148 [](record const & entry, std::string_view k) { return entry.key < k; });
149 if (found != entries.end() && found->key == key) return found->value;
150 }
151 return std::nullopt;
152 }
153
154 std::size_t live_size() const noexcept { return state_->live_size; }
155 element signature() const { return state_->owner.signature(); }
156 std::span<pin const> pins() const noexcept { return state_->owner.pins(); }
157 pin_owner const & owner() const noexcept { return state_->owner; }
158
159 // A snapshot/fork shares this exact immutable state and pin set.
160 reference_world snapshot() const { return *this; }
161
162 std::vector<record> resolved() const {
163 std::map<std::string, std::optional<V>, std::less<>> table;
164 for (auto const & run_pin : state_->owner.pins())
165 for (auto const & entry : run_pin->records()) table[entry.key] = entry.value;
166 std::vector<record> result;
167 result.reserve(state_->live_size);
168 for (auto const & [key, value] : table)
169 if (value) result.push_back({key, value});
170 return result;
171 }
172
174 auto result = A::zero();
175 for (auto const & entry : resolved())
176 result = A::add(result, table_fingerprint<A>::binding(entry.key, entry.value, hash_));
177 return result;
178 }
179
180 // Eager reference compaction. Old snapshots retain their original pins.
182 auto result = from_records(resolved(), hash_);
183 std::vector<std::string> inputs;
184 for (auto const & item : state_->owner.entries()) inputs.push_back(item.object_id);
185 auto replacements = result.owner().entries();
186 auto next = std::make_shared<state>();
187 next->owner = state_->owner.replace(inputs,
188 std::vector<typename pin_owner::entry>(replacements.begin(), replacements.end()));
189 next->live_size = state_->live_size;
190 result.state_ = std::move(next);
191 return result;
192 }
193
194 // Debug resolved-table dump (.rc): every live entry in sorted order, with
195 // full keys without front coding followed by codec values. This is not an
196 // intended access pattern; catalog saves retain object roots separately.
197 // Callers own atomic file replacement, durability and codec/hash agreement.
198 template <class C = u64_world_codec>
199 void export_rc(std::ostream & out, C codec = {}) const {
200 constexpr std::string_view magic{"EVRT.RC\0", 8};
201 out.write(magic.data(), magic.size());
203 world_detail::write_u64(out, C::format_tag);
205 for (auto const & entry : resolved()) {
206 world_detail::write_u64(out, entry.key.size());
207 out.write(entry.key.data(), static_cast<std::streamsize>(entry.key.size()));
208 codec.write(out, *entry.value);
209 }
210 if (!out) throw std::runtime_error("world debug export write failed");
211 }
212
213 // Materializes a fresh reference table from a debug dump, rather than
214 // reopening pinned object roots.
215 template <class C = u64_world_codec>
216 static reference_world import_rc(std::istream & in, C codec = {}, H hash = {},
217 world_import_limits limits = {}) {
218 std::array<char, 8> magic;
219 in.read(magic.data(), magic.size());
220 if (!in || std::string_view(magic.data(), magic.size()) != std::string_view{"EVRT.RC\0", 8})
221 throw std::runtime_error("invalid world debug export header");
222 if (world_detail::read_u64(in) != 1)
223 throw std::runtime_error("unsupported world debug export version");
224 if (world_detail::read_u64(in) != C::format_tag)
225 throw std::runtime_error("world debug export value codec mismatch");
226 auto count = world_detail::read_u64(in);
227 if (count > limits.max_records || count > std::numeric_limits<std::size_t>::max())
228 throw std::runtime_error("world debug export record limit");
229 std::vector<record> entries;
230 entries.reserve(static_cast<std::size_t>(count));
231 std::uint64_t remaining = limits.max_key_bytes;
232 for (std::uint64_t i = 0; i < count; ++i) {
233 auto size = world_detail::read_u64(in);
234 if (size > remaining || size > std::numeric_limits<std::size_t>::max() ||
235 size > std::uint64_t(std::numeric_limits<std::streamsize>::max()))
236 throw std::runtime_error("world debug export key byte limit");
237 remaining -= size;
238 std::string key(static_cast<std::size_t>(size), '\0');
239 in.read(key.data(), static_cast<std::streamsize>(size));
240 if (!in) throw std::runtime_error("truncated world debug export key");
241 if (!entries.empty() && !(entries.back().key < key))
242 throw std::runtime_error("world debug export keys are not strictly sorted");
243 entries.push_back({std::move(key), codec.read(in)});
244 }
245 return from_records(std::move(entries), std::move(hash));
246 }
247
248 private:
249 struct state {
251 std::size_t live_size = 0;
252 };
253
254 reference_world append(std::vector<world_edit<V>> const & edits) const {
255 auto next = std::make_shared<state>(*state_);
256 std::vector<record> records;
257 records.reserve(edits.size());
258 auto contribution = A::zero();
259 auto native_signature = A::zero();
260 for (auto const & edit : edits) {
261 contribution = A::add(contribution,
262 table_fingerprint<A>::delta(edit.key, edit.before, edit.after, hash_));
263 native_signature = A::add(native_signature,
264 table_fingerprint<A>::binding(edit.key, edit.after, hash_));
265 if (!edit.before && edit.after) ++next->live_size;
266 if (edit.before && !edit.after) --next->live_size;
267 records.push_back({edit.key, edit.after});
268 }
269 if (!records.empty()) {
270 auto object = std::make_shared<run>(std::move(records));
271 next->owner = next->owner.add({object->object_id(), object, contribution, native_signature});
272 }
273 reference_world result(*this);
274 result.state_ = std::move(next);
275 return result;
276 }
277
278 template <class, class, class, class> friend struct partition_round;
279 std::shared_ptr<state const> state_;
281 };
282
283 template <class V, class A = wrapping_fingerprint_algebra> struct world_batch {
284 // round_id is supplied by the session protocol and is not derived from the
285 // collision-prone fingerprint. All peers must agree on its starting world.
286 std::string round_id;
287 std::string batch_id;
288 std::uint64_t owner = 0;
289 typename A::element base_signature = A::zero();
290 // Advertised additive contribution of the complete changeset. Receivers
291 // recompute it from the validated old/new bindings before accepting it.
292 typename A::element delta = A::zero();
293 std::vector<world_edit<V>> edits;
294 bool operator==(world_batch const &) const = default;
295 };
296
298
299 // P is a pure, deterministic function string_view -> uint64_t. Partitions
300 // need not be key ranges. Builders may read the same immutable base in
301 // parallel; a receiver serializes apply() calls on its own round accumulator.
302 template <class V, class A, class H, class P> struct partition_round {
305
306 partition_round(world base, std::string round_id, P partition)
307 : base_(std::move(base)), current_(base_), round_id_(std::move(round_id)),
308 partition_(std::move(partition)) {
309 if (round_id_.empty()) throw std::invalid_argument("empty world round identity");
310 }
311
312 batch make_batch(std::string batch_id, std::uint64_t owner,
313 std::vector<world_record<V>> writes) const {
314 if (batch_id.empty()) throw std::invalid_argument("empty world batch identity");
315 std::sort(writes.begin(), writes.end(),
316 [](auto const & a, auto const & b) { return a.key < b.key; });
317 batch result{round_id_, std::move(batch_id), owner, base_.signature(), A::zero(), {}};
318 for (std::size_t i = 0; i < writes.size(); ++i) {
319 auto const & entry = writes[i];
320 if (i && writes[i - 1].key == entry.key)
321 throw std::invalid_argument("duplicate world write key");
322 if (partition_(std::string_view(entry.key)) != owner)
323 throw std::invalid_argument("world write outside assigned partition");
324 auto old = base_.get(entry.key);
325 if (!old && !entry.value) throw std::invalid_argument("deleting absent world key");
326 result.delta = A::add(result.delta,
327 table_fingerprint<A>::delta(entry.key, old, entry.value, base_.hash_));
328 result.edits.push_back({entry.key, std::move(old), entry.value});
329 }
330 return result;
331 }
332
334 if (update.round_id != round_id_ || update.base_signature != base_.signature())
335 throw std::invalid_argument("world batch does not name this round base");
336 if (update.batch_id.empty()) throw std::invalid_argument("empty world batch identity");
337 auto prior = accepted_.find(update.batch_id);
338 if (prior != accepted_.end()) {
339 if (prior->second != update) throw std::invalid_argument("world batch identity reused");
341 }
342 auto delta = A::zero();
343 for (std::size_t i = 0; i < update.edits.size(); ++i) {
344 auto const & edit = update.edits[i];
345 if (i && !(update.edits[i - 1].key < edit.key))
346 throw std::invalid_argument("world batch keys must be strictly sorted");
347 if (partition_(std::string_view(edit.key)) != update.owner)
348 throw std::invalid_argument("world edit outside assigned partition");
349 if (claimed_.contains(edit.key)) throw std::invalid_argument("overlapping world updates");
350 if (base_.get(edit.key) != edit.before)
351 throw std::invalid_argument("world edit old value mismatch");
352 if (!edit.before && !edit.after)
353 throw std::invalid_argument("deleting absent world key");
354 delta = A::add(delta,
355 table_fingerprint<A>::delta(edit.key, edit.before, edit.after, base_.hash_));
356 }
357 if (delta != update.delta)
358 throw std::invalid_argument("world batch contribution mismatch");
359
360 // Stage all allocations before publication. A rejected or failed batch
361 // leaves the current snapshot and its replay/overlap bookkeeping intact.
362 auto next = current_.append(update.edits);
363 auto accepted = accepted_;
364 auto claimed = claimed_;
365 accepted.emplace(update.batch_id, update);
366 for (auto const & edit : update.edits) claimed.emplace(edit.key, update.batch_id);
367 accepted_.swap(accepted);
368 claimed_.swap(claimed);
369 current_ = std::move(next);
371 }
372
373 world const & base() const noexcept { return base_; }
374 world snapshot() const { return current_; }
375 std::size_t accepted_batches() const noexcept { return accepted_.size(); }
376
377 private:
380 std::string round_id_;
382 std::map<std::string, batch, std::less<>> accepted_;
383 std::map<std::string, std::string, std::less<>> claimed_;
384 };
385
386 template <class V, class A, class H, class P>
388}
constexpr std::array< std::byte, 8 > magic
Definition runtime_checkpoint.h:24
bit_string key(key_t< S > const &value)
Definition typed_world.h:89
void write_u64(std::ostream &out, std::uint64_t value)
Definition world.h:69
std::uint64_t read_u64(std::istream &in)
Definition world.h:76
Definition active_engine.h:18
world_apply_result
Definition world.h:297
Declares Everett's pins support.
Definition world.h:302
std::map< std::string, batch, std::less<> > accepted_
Definition world.h:382
world base_
Definition world.h:378
std::map< std::string, std::string, std::less<> > claimed_
Definition world.h:383
P partition_
Definition world.h:381
world_apply_result apply(batch const &update)
Definition world.h:333
partition_round(world base, std::string round_id, P partition)
Definition world.h:306
batch make_batch(std::string batch_id, std::uint64_t owner, std::vector< world_record< V > > writes) const
Definition world.h:312
world snapshot() const
Definition world.h:374
std::string round_id_
Definition world.h:380
std::size_t accepted_batches() const noexcept
Definition world.h:375
world const & base() const noexcept
Definition world.h:373
world current_
Definition world.h:379
Definition world.h:249
std::size_t live_size
Definition world.h:251
pin_owner owner
Definition world.h:250
Definition world.h:110
reference_world append(std::vector< world_edit< V > > const &edits) const
Definition world.h:254
std::shared_ptr< state const > state_
Definition world.h:279
std::size_t live_size() const noexcept
Definition world.h:154
void export_rc(std::ostream &out, C codec={}) const
Definition world.h:199
world_record< V > record
Definition world.h:113
std::optional< V > get(std::string_view key) const
Definition world.h:143
reference_world(H hash={})
Definition world.h:118
element signature() const
Definition world.h:155
std::span< pin const > pins() const noexcept
Definition world.h:156
std::vector< record > resolved() const
Definition world.h:162
reference_world compact() const
Definition world.h:181
pin_owner const & owner() const noexcept
Definition world.h:157
static reference_world from_records(std::vector< record > records, H hash={})
Definition world.h:121
H hash_
Definition world.h:280
V value_type
Definition world.h:111
typename A::element element
Definition world.h:112
static reference_world import_rc(std::istream &in, C codec={}, H hash={}, world_import_limits limits={})
Definition world.h:216
std::shared_ptr< run const > pin
Definition world.h:115
element recompute_signature() const
Definition world.h:173
reference_world snapshot() const
Definition world.h:160
Definition fingerprint.h:56
Definition fingerprint.h:33
Definition world.h:86
static constexpr std::uint64_t format_tag
Definition world.h:89
static void write(std::ostream &out, std::uint64_t value)
Definition world.h:90
static std::uint64_t read(std::istream &in)
Definition world.h:93
Definition world.h:283
std::vector< world_edit< V > > edits
Definition world.h:293
A::element base_signature
Definition world.h:289
bool operator==(world_batch const &) const =default
std::string batch_id
Definition world.h:287
A::element delta
Definition world.h:292
std::uint64_t owner
Definition world.h:288
std::string round_id
Definition world.h:286
Definition world.h:42
bool operator==(world_edit const &) const =default
std::optional< V > before
Definition world.h:44
std::optional< V > after
Definition world.h:45
std::string key
Definition world.h:43
Definition world.h:97
std::uint64_t max_key_bytes
Definition world.h:99
std::uint64_t max_records
Definition world.h:98
Definition world.h:36
std::string key
Definition world.h:37
bool operator==(world_record const &) const =default
std::optional< V > value
Definition world.h:38
Definition world.h:49
std::string const & object_id() const noexcept
Definition world.h:54
std::string object_id_
Definition world.h:65
static std::string next_identity()
Definition world.h:57
std::vector< world_record< V > > records_
Definition world.h:64
world_run(std::vector< world_record< V > > records)
Definition world.h:50
std::span< world_record< V > const > records() const noexcept
Definition world.h:53
Definition fingerprint.h:22