Everett
Loading...
Searching...
No Matches
replacement_rebuild.h
Go to the documentation of this file.
1
12#pragma once
13
15#include <everett/typed_scan.h>
16#include <deque>
17#include <unordered_map>
18
19namespace everett {
20 namespace replacement_detail {
21 template <class Family, class = void> struct clean_family {
22 using type = Family;
23 static constexpr bool enabled = false;
24 };
25 template <class Family> struct clean_family<Family, std::void_t<typename Family::storage_type::clean_storage_type>> {
26 using type = typename Family::template rebind_storage<typename Family::storage_type::clean_storage_type>;
27 static constexpr bool enabled = !std::is_same_v<Family, type>;
28 };
29 }
30 // Trusted semantic checkpoint extension; the runtime frontier is stored
31 // separately. Generation fields never authenticate the payload or its hash.
32 template <class A = wrapping_fingerprint_algebra>
35 std::uint64_t clean_base = 0, mutations = 0;
36 bool rebuilding = false;
38 replacement_metadata(base_type value, std::uint64_t b, std::uint64_t u, bool active)
39 : base_type(std::move(value)), clean_base(b), mutations(u), rebuilding(active) {}
40 void validate(std::uint64_t admissions) const {
41 if (clean_base > admissions || mutations != admissions - clean_base ||
42 this->live_count > admissions ||
43 (clean_base > mutations && this->live_count < clean_base - mutations))
44 throw std::invalid_argument("invalid replacement generation mass");
45 auto trigger = clean_base / 4;
46 if (!rebuilding) {
47 if ((clean_base < 64 && (mutations || this->live_count != clean_base)) ||
48 (clean_base >= 64 && mutations >= trigger))
49 throw std::invalid_argument("inactive replacement generation passed its trigger");
50 } else {
51 // At freeze n_s <= b + floor(b/4); at most floor(n_s/8)
52 // intervening admissions fit before handoff. Divide before adding to
53 // avoid overflow in the upper bound for extreme metadata.
54 auto extra = clean_base / 8 + trigger / 8 + (clean_base % 8 + trigger % 8) / 8;
55 if (clean_base < 64 || mutations < trigger || mutations - trigger > extra)
56 throw std::invalid_argument("invalid active replacement generation");
57 }
58 }
59 std::vector<std::byte> encode() const requires std::same_as<typename A::element, std::uint64_t> {
60 auto inner = base_type::encode();
61 std::vector<std::byte> result(40 + inner.size());
62 constexpr std::array<unsigned char, 8> magic{'E', 'V', 'R', 'T','.','R','B',0};
63 for (unsigned i = 0; i != 8; ++i) result[i] = std::byte(magic[i]);
64 std::array<std::uint64_t, 4> fields{1, clean_base, mutations, rebuilding ? 1u : 0u};
65 for (unsigned n = 0; n != fields.size(); ++n)
66 for (unsigned i = 0; i != 8; ++i) result[8 + n * 8 + i] = std::byte(fields[n] >> (i << 3));
67 std::copy(inner.begin(), inner.end(), result.begin() + 40);
68 return result;
69 }
70 static replacement_metadata decode(std::span<std::byte const> data)
71 requires std::same_as<typename A::element, std::uint64_t> {
72 constexpr std::array<unsigned char, 8> magic{'E', 'V', 'R', 'T','.','R','B',0};
73 if (data.size() <= 56) throw std::invalid_argument("truncated replacement metadata");
74 for (unsigned i = 0; i != 8; ++i)
75 if (data[i] != std::byte(magic[i])) throw std::invalid_argument("replacement metadata signature mismatch");
76 std::array<std::uint64_t, 4> fields{};
77 for (unsigned n = 0; n != fields.size(); ++n)
78 for (unsigned i = 0; i != 8; ++i)
79 fields[n] |= std::uint64_t(std::to_integer<unsigned char>(data[8 + n * 8 + i])) << (i << 3);
80 if (fields[0] != 1 || fields[3] > 1) throw std::invalid_argument("replacement metadata version or flags");
81 return {base_type::decode(data.subspan(40)), fields[1], fields[2], fields[3] != 0};
82 }
83 bool operator==(replacement_metadata const &) const = default;
84 };
85
86 template <class P = storage_policy<>, class A = wrapping_fingerprint_algebra,
87 class Family = redundant_runtime_family<P>>
88 struct replacement_world : typed_world<P, A, Family> {
91 using runtime_snapshot = typename Family::snapshot_type;
92 replacement_world(base_type value, std::uint64_t b = 0, std::uint64_t u = 0, bool active = false)
93 : base_type(std::move(value)), metadata_(std::make_shared<metadata_type const>(base_type::metadata(), b, u, active)) {
94 metadata_->validate(this->runtime().admissions());
95 }
96 metadata_type const & metadata() const & noexcept { return *metadata_; }
97 static replacement_world restore(runtime_snapshot data, metadata_type metadata, std::string_view schema) {
98 metadata.validate(data.admissions());
100 return {base_type::restore(std::move(data), std::move(metadata), schema), b, u, active};
101 }
102 private:
103 std::shared_ptr<metadata_type const> metadata_;
104 };
105
107 std::uint64_t mutations = 0, generations = 0;
108 std::uint64_t foreground_charged = 0, candidate_charged = 0;
109 std::uint64_t reserved = 0, granted = 0, committed = 0;
110 std::uint64_t scan_records = 0, clean_rows = 0, replayed = 0;
113 };
115 std::uint64_t clean_base = 0, mutations = 0;
116 std::uint64_t frozen_live = 0, horizon = 0, admitted = 0, replayed = 0;
117 std::uint64_t initial_bound = 0, quantum = 0, action_bound = 0;
118 std::uint64_t committed = 0, credit = 0, queued = 0;
119 std::uint64_t source_records = 0, clean_rows = 0;
120 bool rebuilding = false, scanning = false;
121 };
122
123 // This initial executor supports one occupied replacement sort. A custom
124 // sort supplies clean(key,state)->arrow; the default optional-string sort
125 // already represents a clean state with the same replacement arrow.
126 // Budgets are structural allowances, not byte counts or elapsed time.
127 template <class P = storage_policy<>, class A = wrapping_fingerprint_algebra,
128 std::uint64_t DepthLimit = 256, class Family = redundant_runtime_family<P>>
130 using policy_type = P;
132 static_assert(!std::is_void_v<sort_type>, "replacement rebuild needs one occupied sort");
137 static_assert(typed_detail::replacement<sort_type> && std::is_same_v<state_type, arrow_type>,
138 "replacement rebuild requires replacement state/arrow types");
142 static_assert(charged_service, "replacement rebuild requires charged redundant service");
143 using runtime_family = Family;
149 static_assert(DepthLimit >= 3);
150 static constexpr std::uint64_t small_limit = 64;
151 static constexpr std::uint64_t tiny_record_limit = 256;
152
154 explicit replacement_rebuild_engine(std::string schema)
155 : foreground_(std::make_unique<engine_type>(std::move(schema))), published_(foreground_->snapshot()) { work_.foreground_charged = foreground_->work().charged; }
157 if (source.runtime().admissions() != source.metadata().live_count)
158 throw std::invalid_argument("replacement rebuild restore needs a clean admission mass");
159 auto b = source.metadata().live_count;
160 return from_snapshot(world_type(std::move(source), b));
161 }
163 return replacement_rebuild_engine(std::move(source));
164 }
165 template <class Storage>
167 requires requires { engine_type::from_snapshot(source, std::move(storage)); }
168 {
169 if (source.runtime().admissions() != source.metadata().live_count)
170 throw std::invalid_argument("replacement rebuild restore needs a clean admission mass");
171 auto b = source.metadata().live_count;
172 return from_snapshot(world_type(std::move(source), b), std::move(storage));
173 }
174 template <class Storage>
176 requires requires { engine_type::from_snapshot(source, std::move(storage)); }
177 {
178 return replacement_rebuild_engine(std::move(source), std::move(storage));
179 }
183 replacement_rebuild_engine & operator=(replacement_rebuild_engine &&) noexcept = default;
184
185 world_type snapshot() const { active(); return published_; }
186 bool failed() const noexcept {
187 return failed_ || (foreground_ && foreground_->failed()) ||
188 (job_ && job_->candidate && job_->candidate->failed());
189 }
190 // A shared storage context is serialized with both private executors.
191 // Poison it too when an outer durable publication becomes uncertain.
192 void poison() noexcept {
193 failed_ = true;
194 if (foreground_) foreground_->poison();
195 if (job_ && job_->candidate) job_->candidate->poison();
196 }
197 auto storage() const requires requires(engine_type const & core) { core.storage(); } {
198 active();
199 return foreground_->storage();
200 }
201 // Replace only a settled, equivalent physical graph. The concrete storage
202 // context survives publication, including after a candidate handoff.
203 void rebase(world_type state) {
204 writable();
205 if (pending() || state.metadata() != published_.metadata() ||
206 state.runtime().admissions() != published_.runtime().admissions())
207 throw std::invalid_argument("replacement rebase requires a settled equivalent snapshot");
208 try {
209 foreground_->rebase(state);
211 published_ = std::move(state);
212 } catch (...) { poison(); throw; }
213 }
214 bool pending() const noexcept { return foreground_ && (job_ || foreground_->pending()); }
215 bool admission_ready() const noexcept { return foreground_ && !failed() && !recovering_ && foreground_->admission_ready(); }
216 replacement_rebuild_work work() const noexcept { return work_; }
219 if (job_) {
220 auto const & j = *job_; out.frozen_live = j.live; out.horizon = j.horizon;
221 out.admitted = j.admitted; out.replayed = j.replayed; out.initial_bound = j.bound;
222 out.quantum = j.quantum; out.action_bound = j.action; out.committed = j.committed;
223 out.credit = j.credit; out.queued = j.queue.size(); out.source_records = j.source_records; out.clean_rows = j.rows; out.rebuilding = true; out.scanning = j.building;
224 }
225 return out;
226 }
227 static auto batch() { return engine_type::batch(); }
228 template <class S = sort_type> requires std::same_as<S, sort_type>
229 static contribution_type put(key_type const & key, state_type const & value) {
230 return engine_type::template put<sort_type>(key, value);
231 }
232 template <class S = sort_type> requires std::same_as<S, sort_type>
233 static contribution_type erase(key_type const & key) { return engine_type::template erase<sort_type>(key); }
234
235 template <class S = sort_type> requires std::same_as<S, sort_type>
236 static contribution_type change(key_type const & key, arrow_type const & arrow) {
237 return engine_type::template change<S>(key, arrow);
238 }
239 // Accepted input allowance only. Recovery of already published history is
240 // separately serviced behind admission_ready(), before a session claims input.
241 static std::uint64_t reservation_work(std::uint64_t records) {
242 auto h = std::min<std::uint64_t>(64, DepthLimit - 3);
243 auto g = action_bound(h), c = runtime_type::local_charge_bound;
244 constexpr std::uint64_t runs = 128, scan = 32, setup = runs * (scan + 8) + 32;
245 auto large = add(add(mul(22, g), mul(mul(10, add(c, 32)), add(h, 1))), setup + 26 * scan + 2);
246 // Below the small threshold the valid generation has at most 256
247 // physical occurrences. L <= 64 and bit_width(L) <= 7.
248 auto small = add(setup + 321 * scan, add(mul(130, action_bound(7)), mul(512, add(c, 32))));
249 auto query = mul(mul(64, add(DepthLimit, 1)), add(add(P::group_size, P::codec_block_size), 16));
250 auto extra = add(add(std::max(large, small), query), runs * 8 + 32);
252 return add(engine_type::reservation_work(records), mul(records, extra));
253 }
255 auto quote = engine_type::reservation(input);
256 quote.work = reservation_work(input.records().size());
257 return quote;
258 }
259
260 // Complete validation precedes mutation. The per-key fallback also prepares
261 // detached replay records before execution. Intermediate cuts stay private
262 // until the entire batch and its owed service succeeds.
264 writable();
265 if (!admission_ready()) throw std::logic_error("replacement foreground needs recovery service");
266 if (input.base() && input.base()->metadata().schema_id != published_.metadata().schema_id)
267 throw std::invalid_argument("rebuild contribution uses another schema");
268 if (published_.runtime().query_root().head()->depth() > DepthLimit ||
269 (input.base() && input.base()->runtime().query_root().head()->depth() > DepthLimit))
270 throw std::length_error("replacement query exceeds depth allowance");
271 if (input.records().empty()) return published_;
272 // Initial unique string replacements are already clean arrows. Share
273 // the typed preflight and avoid constructing detached per-key replay
274 // entries when the runtime can initialize a prefix and admit the tail
275 // before publishing the complete pristine batch.
276 if constexpr (std::is_same_v<sort_type, unsorted<std::optional<std::string>>> &&
277 requires(runtime_type & runtime, std::span<profile_record const> records) {
278 runtime.try_initialize_sorted(records, std::uint64_t{}, DepthLimit);
279 }) {
280 auto count = input.records().size();
281 if (count >= 2 && !base_ && !mutations_ &&
282 !job_ && !recovering_ && !foreground_->pending() && !mass(published_) &&
284 auto metadata = foreground_->prepare(input);
285 auto accepted = add(work_.mutations, count);
286 auto prior = foreground_->work().charged;
287 try {
288 if (foreground_->initialize(input, metadata)) {
290 prior = foreground_->work().charged;
291 base_ = count;
292 work_.mutations = accepted;
294 return published_;
295 }
296 } catch (...) {
298 poison();
299 throw;
300 }
301 }
302 }
303 std::vector<mutation> entries; entries.reserve(input.records().size());
304 foreground_->visit_changes(input, [&]<class S>(std::type_identity<S>, auto const & key,
305 auto && before, auto && after, auto const & record) {
306 static_assert(std::is_same_v<S, sort_type>);
307 entries.push_back({key, std::move(before), std::move(after), typed_detail::value<P, S>(record.value.view()),
308 0, record.retained_limit_bits});
309 });
310 try {
311 auto generation = work_.generations;
312 for (auto & entry : entries) {
313 // Ordinary carries preserve all distinct keys. A clean-generation
314 // handoff can remove a predecessor and require a longer literal.
315 if (!semantics::present(entry.key, entry.after) && work_.generations != generation)
316 entry.retained_limit_bits = 0;
317 apply(std::move(entry));
318 }
320 return published_;
321 } catch (...) { poison(); throw; }
322 }
323 std::optional<world_type> advance(std::uint64_t budget) {
324 writable(); if (!budget || !pending()) return {};
325 auto before = published_;
326 try {
327 if (job_ && foreground_->admission_ready()) { grant(budget); service(); }
328 else foreground_advance(budget);
330 } catch (...) { poison(); throw; }
331 if (before.runtime().same_layout(published_.runtime()) && before.metadata() == published_.metadata()) return {};
332 return published_;
333 }
334
335 private:
336 struct mutation {
340 std::uint64_t ordinal;
341 std::optional<std::uint64_t> retained_limit_bits;
342 };
343 static contribution_type command(mutation const & entry) {
344 auto result = engine_type::template change<sort_type>(entry.key, entry.arrow);
345 result.records_[0].retained_limit_bits = entry.retained_limit_bits;
346 return result;
347 }
348 struct rebuild {
350 std::unique_ptr<scan_type> scan;
351 std::unique_ptr<engine_type> candidate;
352 std::deque<mutation> queue;
353 std::uint64_t live = 0, horizon = 1, frozen_ordinal = 0;
354 std::uint64_t admitted = 0, replayed = 0, rows = 0, source_records = 0;
355 std::uint64_t bound = 0, quantum = 0, action = 0, scan_price = 0, setup = 0;
356 std::uint64_t committed = 0, credit = 0;
357 bool building = true;
358 bool tiny = false;
359 explicit rebuild(typed_world_type value) : frozen(std::move(value)) {}
360 };
361 std::unique_ptr<engine_type> foreground_;
363 std::unique_ptr<rebuild> job_;
364 std::uint64_t base_ = 0, mutations_ = 0;
366 bool failed_ = false, recovering_ = false;
367
369 : foreground_(std::make_unique<engine_type>(engine_type::from_snapshot(value))), published_(std::move(value)),
370 base_(published_.metadata().clean_base), mutations_(published_.metadata().mutations) {
372 }
373 template <class Storage>
375 : foreground_(std::make_unique<engine_type>(engine_type::from_snapshot(value, std::move(storage)))),
376 published_(std::move(value)), base_(published_.metadata().clean_base),
377 mutations_(published_.metadata().mutations) {
379 }
381 if (published_.runtime().query_root().head()->depth() > DepthLimit)
382 throw std::length_error("restored replacement query exceeds depth allowance");
383 work_.foreground_charged = foreground_->work().charged;
384 if (published_.metadata().rebuilding) { recovering_ = true; start(true); }
385 }
386 world_type publication() const { return {foreground_->snapshot(), base_, mutations_, bool(job_)}; }
387 void active() const { if (!foreground_) throw std::logic_error("moved-from replacement rebuild engine"); }
388 void writable() const { active(); if (failed()) throw std::logic_error("failed replacement rebuild engine"); }
389 static std::uint64_t mass(typed_world_type const & state) { return state.runtime().admissions(); }
390 static std::uint64_t add(std::uint64_t a, std::uint64_t b) { return profile_detail::add(a, b); }
391 static std::uint64_t mul(std::uint64_t a, std::uint64_t b) { return profile_detail::multiply(a, b); }
392 static std::uint64_t ceil(std::uint64_t a, std::uint64_t b) { return a / b + (a % b != 0); }
393 static void require(bool value, char const * message) { if (!value) throw std::logic_error(message); }
394 static arrow_type clean(key_type const & key, state_type const & state) {
395 if constexpr (requires { semantics::clean(key, state); }) return semantics::clean(key, state);
396 else {
397 static_assert(std::is_same_v<sort_type, unsorted<std::optional<std::string>>>,
398 "custom replacement sort must define clean(key,state)");
399 return state;
400 }
401 }
402 void foreground_advance(std::uint64_t budget) {
403 auto prior = foreground_->work().charged;
404 try { foreground_->advance(budget); }
405 catch (...) { work_.foreground_charged = add(work_.foreground_charged, foreground_->work().charged - prior); throw; }
407 }
408 void apply(mutation entry) {
409 require(foreground_->admission_ready(), "replacement foreground exhausted its admission service");
410 entry.ordinal = add(work_.mutations, 1);
411 // A fresh replacement in a clean string table adds one live row and
412 // one physical admission, without introducing history to collect.
413 // Extend that clean base; pending ordinary carries still receive their
414 // normal admission service. Once history exists, every write funds it.
415 bool extends_clean = false;
416 if constexpr (std::is_same_v<sort_type, unsorted<std::optional<std::string>>>)
417 extends_clean = !job_ && !mutations_ && !entry.before && bool(entry.after);
418 if (job_) {
419 require(job_->admitted < job_->horizon, "rebuild deadline exhausted before admission");
420 job_->queue.push_back(entry);
421 }
422 auto prior = foreground_->work().charged;
423 try {
424 auto input = command(entry);
425 foreground_->template contribute_validated<sort_type>(input.records_[0], entry.key, entry.before, entry.after);
426 }
427 catch (...) { work_.foreground_charged = add(work_.foreground_charged, foreground_->work().charged - prior); throw; }
429 work_.mutations = entry.ordinal;
430 if (extends_clean) {
431 base_ = add(base_, 1);
432 require(mass(foreground_->snapshot()) == base_, "extended clean generation mass mismatch");
433 return;
434 }
436 require(mass(foreground_->snapshot()) == add(base_, mutations_), "foreground generation mass mismatch");
437 if (job_) {
438 ++job_->admitted;
439 work_.reserved = add(work_.reserved, job_->action);
440 work_.maximum_replay = std::max<std::uint64_t>(work_.maximum_replay, job_->queue.size());
441 grant(add(job_->quantum, job_->action)); service();
442 require(!job_ || job_->admitted < job_->horizon, "rebuild missed funded handoff horizon");
443 } else {
444 auto live = foreground_->snapshot().metadata().live_count;
445 bool small = base_ < small_limit || live < small_limit;
446 if (small || mutations_ >= base_ / 4) {
447 start(small);
448 if (small) { grant(job_->bound); service(); require(!job_, "small rebuild did not finish"); }
449 }
450 }
451 }
452 static std::uint64_t action_bound(std::uint64_t height) {
453 auto depth = add(height, 3), c = runtime_type::local_charge_bound;
454 auto ready = add(add(mul(2, c), mul(16, depth)), 512);
455 auto query = mul(mul(64, add(depth, 1)), add(add(P::group_size, P::codec_block_size), 16));
456 auto service = mul(mul(8, c), add(height, 2));
457 return add(add(ready, service), add(query, add(mul(8, height), 16)));
458 }
459 void start(bool small) {
460 auto source = foreground_->snapshot();
461 auto freeze = add(mul(source.runtime().runs().size(), 8), 32);
462 work_.reserved = add(work_.reserved, freeze);
463 work_.granted = add(work_.granted, freeze);
464 work_.committed = add(work_.committed, freeze);
465 auto next = std::make_unique<rebuild>(std::move(source));
466 next->live = next->frozen.metadata().live_count;
467 next->horizon = small ? 1 : next->live / 8;
468 require(next->horizon, "empty rebuilding horizon");
469 next->frozen_ordinal = work_.mutations;
470 auto limit = add(next->live, small ? 0 : next->horizon);
471 auto height = std::bit_width(limit); auto depth = add(height, 3);
472 if (depth > DepthLimit) throw std::length_error("rebuild candidate exceeds supported depth");
473 auto c = runtime_type::local_charge_bound;
474 next->action = action_bound(height);
475 std::uint64_t physical = 0;
476 auto runs = next->frozen.runtime().runs();
477 for (auto const & run : runs) physical = add(physical, run->native->size());
478 next->source_records = physical;
480 next->tiny = small && next->live <= small_limit && physical <= tiny_record_limit;
481 next->scan_price = add(mul(2, std::bit_width(runs.size())), 16);
482 next->setup = add(mul(runs.size(), add(next->scan_price, 8)), 32);
483 auto r = add(next->setup, mul(add(add(physical, next->live), 1), next->scan_price));
484 r = add(r, mul(next->live, next->action));
485 // All maintenance through the largest possible replay cut is reserved,
486 // including a large carry first triggered by an intervening mutation.
487 r = add(r, mul(mul(add(c, 32), limit), add(height, 1)));
488 r = add(r, mul(limit, next->action)); // Nested grant fragments at settlement.
489 r = add(r, mul(next->horizon, next->action)); // Outer atomic-action fragments.
490 next->bound = add(r, next->action); // Final metadata check and ownership handoff.
491 if (next->tiny) next->bound = add(next->bound, tiny_conversion_bound());
492 next->quantum = ceil(next->bound, next->horizon);
493 work_.reserved = add(work_.reserved, next->bound);
494 job_ = std::move(next);
495 }
496 void grant(std::uint64_t amount) {
497 auto credit = add(job_->credit, amount), total = add(work_.granted, amount);
498 job_->credit = credit; work_.granted = total;
499 }
500 std::uint64_t price() const {
501 auto const & j = *job_;
502 if (j.tiny) return j.bound;
503 if (!j.candidate) return j.setup;
504 if (j.building && !j.scan->has_row() && !j.scan->done()) return j.scan_price;
505 return j.action;
506 }
507 void candidate_work(auto && operation) {
508 auto before = job_->candidate->work().charged;
509 try { operation(); }
510 catch (...) {
511 work_.candidate_charged = add(work_.candidate_charged, job_->candidate->work().charged - before);
512 throw;
513 }
514 work_.candidate_charged = add(work_.candidate_charged, job_->candidate->work().charged - before);
515 }
516 // A settled candidate with at most 64 admissions has at most seven levels.
517 // Checked restore confines every object route to the three slots per
518 // level. Their pairs, at most one carrier per level, and the prepared
519 // head chain use fewer than 8*(height+1) distinct pairs. Each pair has at most 2*64+4 occurrences:
520 // V <= 64 + ceil(V/K) + ceil(64/K), with K >= 3. Charge an index allowance
521 // including both directory streams, plus object/level admission metadata.
522 static std::uint64_t tiny_conversion_bound() {
523 auto pairs = mul(8, add(std::bit_width(small_limit), 1));
524 return add(mul(pairs, add(32, mul(add(mul(2, small_limit), 4), add(P::group_size, 8)))), 1024);
525 }
526 template <class Source> auto convert_tiny(Source const & source) {
527 using source_pair = typename Source::query_type::pair_type;
528 using source_node = std::remove_const_t<typename source_pair::element_type>;
529 using source_object = typename Source::object_pointer::element_type;
530 using target_node = typename Family::node_type;
531 using target_pair = typename target_node::pair_type;
532 using target_object = typename Family::snapshot_type::object_pointer;
533 static_assert(std::is_same_v<typename source_node::native_type, typename Family::native_type>);
534 auto const & old = source.frontier();
535 require(old.admissions <= small_limit && old.levels.size() <= std::bit_width(small_limit) && !old.service_due,
536 "tiny conversion exceeds its bounded settled frontier");
537 std::uint64_t charged = 0;
538 auto charge = [&](std::uint64_t amount) {
539 charged = add(charged, amount);
540 require(charged <= tiny_conversion_bound(), "tiny conversion exceeded its allowance");
543 };
544 charge(add(32, mul(old.levels.size(), 16)));
545 std::unordered_map<source_node const *, target_pair> pairs;
546 auto copy_pair = [&](auto && self, source_pair const & value) -> target_pair {
547 if (!value) return {};
548 if (auto found = pairs.find(value.get()); found != pairs.end()) return found->second;
549 auto main = self(self, value->main_target());
550 require(pairs.size() < 8 * (std::bit_width(small_limit) + 1), "too many tiny conversion pairs");
551 auto secondary = value->secondary_target();
552 auto a = main ? ceil(main->virtual_size(), P::group_size) : 0;
553 auto b = secondary ? ceil(secondary->size(), P::group_size) : 0;
554 auto n = add(value->native_owner()->size(), add(a, b));
555 require(n <= 2 * small_limit + 4, "tiny pair exceeds its occurrence allowance");
556 charge(add(add(16, mul(n, P::group_size + 6)), add(ceil(a, P::codec_block_size), ceil(b, P::codec_block_size))));
558 while (!builder.done()) builder.step(1);
559 auto result = target_node::from_built(builder.finish());
561 pairs.emplace(value.get(), result); return result;
562 };
563 std::unordered_map<source_object const *, target_object> objects;
564 auto copy_object = [&](auto && self, typename Source::object_pointer const & value) -> target_object {
565 if (!value) return {};
566 if (auto found = objects.find(value.get()); found != objects.end()) return found->second;
567 typename Family::routes_type next{self(self, value->next.main), self(self, value->next.secondary)};
568 require(objects.size() < 3 * old.levels.size(), "too many tiny conversion objects");
569 charge(16);
570 auto result = std::make_shared<typename Family::object_type const>(typename Family::object_type{
571 value->identity, value->first, value->last, value->level, value->native,
572 copy_pair(copy_pair, value->pair), std::move(next)});
573 objects.emplace(value.get(), result); return result;
574 };
575 auto copy_route = [&](auto const & route) {
576 return typename Family::routes_type{copy_object(copy_object, route.main), copy_object(copy_object, route.secondary)};
577 };
578 typename Family::frontier_type out;
579 out.admissions = old.admissions; out.next_identity = old.next_identity; out.service_due = old.service_due;
580 out.levels.resize(old.levels.size());
581 for (std::size_t i = 0; i != old.levels.size(); ++i) {
582 auto const & before = old.levels[i]; auto & after = out.levels[i];
583 require(!before.job, "tiny conversion requires completed jobs");
584 unsigned active = 0;
585 for (std::size_t j = 0; j != before.slots.size(); ++j) {
586 auto const & slot = before.slots[j]; active += slot.state == redundant_slot_state::active;
587 after.slots[j] = {slot.state, copy_object(copy_object, slot.object), copy_route(slot.route),
588 copy_pair(copy_pair, slot.carrier), slot.ever_visible};
589 }
590 require(active < 2, "tiny conversion requires a settled level");
591 after.last_destination = before.last_destination; after.last_destination_visible = before.last_destination_visible;
592 }
593 out.root = copy_route(old.root);
594 auto head = copy_pair(copy_pair, source.query_root().head());
595 auto result = Family::snapshot_type::restore(std::move(out), std::move(head));
596 return result;
597 }
598 void tiny_rebuild() requires replacement_detail::clean_family<Family>::enabled {
599 using clean_family = typename replacement_detail::clean_family<Family>::type;
601 auto & j = *job_;
602 require(j.tiny && j.live <= small_limit && j.source_records <= tiny_record_limit &&
603 j.queue.empty() && !j.admitted, "tiny rebuild escaped its bounded eager path");
604 clean_engine candidate(j.frozen.metadata().schema_id);
605 work_.candidate_charged = add(work_.candidate_charged, candidate.work().charged);
606 auto execute = [&](auto && operation) {
607 auto before = candidate.work().charged;
608 try { operation(); }
609 catch (...) { work_.candidate_charged = add(work_.candidate_charged, candidate.work().charged - before); throw; }
610 work_.candidate_charged = add(work_.candidate_charged, candidate.work().charged - before);
611 };
612 scan_type rows(j.frozen);
613 while (!rows.done()) {
614 if (!rows.has_row()) work_.scan_records = add(work_.scan_records, rows.step(1));
615 if (!rows.has_row()) continue;
616 auto row = rows.take_row();
617 auto arrow = clean(row.key, row.value);
618 require(semantics::apply(row.key, semantics::initial(row.key), arrow) == row.value, "invalid clean replacement arrow");
619 while (!candidate.admission_ready()) execute([&] { candidate.advance(j.action); });
620 execute([&] { candidate.contribute(clean_engine::template change<sort_type>(row.key, arrow)); });
621 ++j.rows; work_.clean_rows = add(work_.clean_rows, 1);
622 require(j.rows <= j.live, "tiny scan exceeds frozen live count");
623 }
624 auto scanned = candidate.snapshot();
625 require(rows.consumed() == j.source_records && j.rows == j.live &&
626 scanned.runtime().admissions() == j.live && scanned.metadata() == j.frozen.metadata(),
627 "tiny rebuild differs from frozen source");
628 while (candidate.pending()) execute([&] { candidate.advance(j.action); });
629 auto clean = candidate.snapshot();
630 auto runtime = convert_tiny(clean.runtime());
631 auto state = typed_world_type::restore(std::move(runtime), clean.metadata(), j.frozen.metadata().schema_id);
632 require(state.metadata() == foreground_->snapshot().metadata() && state.runtime().admissions() == j.live,
633 "tiny handoff differs from foreground");
634 auto replacement = engine_type::from_snapshot(std::move(state), foreground_->storage());
636 base_ = j.live; mutations_ = 0;
638 foreground_ = std::make_unique<engine_type>(std::move(replacement)); job_.reset(); recovering_ = false;
639 }
640 void service() {
641 while (job_) {
642 auto cost = price(); if (job_->credit < cost) break;
643 auto limit = add(job_->bound, mul(job_->admitted, job_->action));
644 require(job_->committed <= limit && cost <= limit - job_->committed, "rebuild exceeded its reserved work bound");
645 job_->credit -= cost; job_->committed += cost; work_.committed = add(work_.committed, cost);
646 auto & j = *job_;
648 if (j.tiny) { tiny_rebuild(); continue; }
649 }
650 if (!j.candidate) {
651 auto seed = std::make_unique<engine_type>(j.frozen.metadata().schema_id);
652 work_.candidate_charged = add(work_.candidate_charged, seed->work().charged);
653 if constexpr (requires { foreground_->storage(); }) {
654 // The empty seed does not execute native work. Its initialization
655 // and context-aware restoration both fit the setup's 32-unit margin.
656 j.candidate = std::make_unique<engine_type>(
657 engine_type::from_snapshot(seed->snapshot(), foreground_->storage()));
658 work_.candidate_charged = add(work_.candidate_charged, j.candidate->work().charged);
659 } else j.candidate = std::move(seed);
660 j.scan = std::make_unique<scan_type>(j.frozen);
661 } else if (j.building) {
662 if (j.scan->has_row()) {
663 if (!j.candidate->admission_ready()) candidate_work([&]{ j.candidate->advance(j.action); });
664 else {
665 auto row = j.scan->take_row();
666 auto arrow = clean(row.key, row.value);
667 require(semantics::apply(row.key, semantics::initial(row.key), arrow) == row.value, "invalid clean replacement arrow");
668 auto input = engine_type::template change<sort_type>(row.key, arrow);
669 candidate_work([&]{ j.candidate->template contribute_validated<sort_type>(
670 input.records_[0], row.key, semantics::initial(row.key), row.value); });
671 ++j.rows; work_.clean_rows = add(work_.clean_rows, 1);
672 }
673 } else if (!j.scan->done()) work_.scan_records = add(work_.scan_records, j.scan->step(1));
674 else {
675 require(j.scan->consumed() == j.source_records && j.rows == j.live && mass(j.candidate->snapshot()) == j.live &&
676 j.candidate->snapshot().metadata() == j.frozen.metadata(), "clean rebuild differs from frozen source");
677 j.scan.reset(); j.building = false;
678 }
679 } else if (!j.queue.empty()) {
680 if (!j.candidate->admission_ready()) candidate_work([&]{ j.candidate->advance(j.action); });
681 else {
682 auto const & entry = j.queue.front();
683 require(entry.ordinal == add(add(j.frozen_ordinal, j.replayed), 1), "rebuild replay order changed");
684 auto input = command(entry);
685 input.records_[0].retained_limit_bits = 0;
686 candidate_work([&]{ j.candidate->template contribute_validated<sort_type>(
687 input.records_[0], entry.key, entry.before, entry.after); });
688 j.queue.pop_front(); ++j.replayed; work_.replayed = add(work_.replayed, 1);
689 }
690 } else if (j.candidate->pending()) candidate_work([&]{ j.candidate->advance(j.action); });
691 else {
692 auto state = j.candidate->snapshot();
693 require(j.replayed == j.admitted && state.runtime().admissions() == add(j.live, j.replayed) &&
694 state.metadata() == foreground_->snapshot().metadata(), "rebuild handoff differs from foreground");
695 base_ = j.live; mutations_ = j.replayed;
698 foreground_ = std::move(j.candidate); job_.reset(); recovering_ = false;
699 }
700 }
701 }
702 };
703}
std::uint64_t add(std::uint64_t a, std::uint64_t b)
Definition profile.h:39
std::uint64_t multiply(std::uint64_t a, std::uint64_t b)
Definition profile.h:44
typename sort_codec< S >::value_codec::value_type arrow_t
Definition typed_world.h:53
typename sort_codec< S >::key_codec::value_type key_t
Definition typed_world.h:52
typename default_sort< typename registry_detail::info< typename P::registry_type >::leaves >::type default_sort_t
Definition typed_world.h:51
typename sort_semantics< S >::state_type state_t
Definition typed_world.h:54
Definition active_engine.h:18
Executes redundant COLA slots with charged jobs and immutable frontiers.
Definition cola_index.h:496
auto finish()
Definition cola_index.h:581
bool done() const noexcept
Definition cola_index.h:520
std::uint64_t step(std::uint64_t budget)
Definition cola_index.h:529
typename Family::template rebind_storage< typename Family::storage_type::clean_storage_type > type
Definition replacement_rebuild.h:26
Definition replacement_rebuild.h:21
static constexpr bool enabled
Definition replacement_rebuild.h:23
Family type
Definition replacement_rebuild.h:22
Definition replacement_rebuild.h:33
std::uint64_t clean_base
Definition replacement_rebuild.h:35
bool rebuilding
Definition replacement_rebuild.h:36
std::uint64_t mutations
Definition replacement_rebuild.h:35
replacement_metadata(base_type value, std::uint64_t b, std::uint64_t u, bool active)
Definition replacement_rebuild.h:38
void validate(std::uint64_t admissions) const
Definition replacement_rebuild.h:40
std::vector< std::byte > encode() const
Definition replacement_rebuild.h:59
static replacement_metadata decode(std::span< std::byte const > data)
Definition replacement_rebuild.h:70
bool operator==(replacement_metadata const &) const =default
Definition replacement_rebuild.h:336
std::uint64_t ordinal
Definition replacement_rebuild.h:340
arrow_type arrow
Definition replacement_rebuild.h:339
std::optional< std::uint64_t > retained_limit_bits
Definition replacement_rebuild.h:341
state_type after
Definition replacement_rebuild.h:338
state_type before
Definition replacement_rebuild.h:338
key_type key
Definition replacement_rebuild.h:337
Definition replacement_rebuild.h:348
std::uint64_t frozen_ordinal
Definition replacement_rebuild.h:353
std::uint64_t source_records
Definition replacement_rebuild.h:354
std::uint64_t action
Definition replacement_rebuild.h:355
std::unique_ptr< scan_type > scan
Definition replacement_rebuild.h:350
std::uint64_t rows
Definition replacement_rebuild.h:354
std::unique_ptr< engine_type > candidate
Definition replacement_rebuild.h:351
std::uint64_t horizon
Definition replacement_rebuild.h:353
std::uint64_t live
Definition replacement_rebuild.h:353
typed_world_type frozen
Definition replacement_rebuild.h:349
std::uint64_t replayed
Definition replacement_rebuild.h:354
std::uint64_t committed
Definition replacement_rebuild.h:356
std::uint64_t admitted
Definition replacement_rebuild.h:354
std::uint64_t credit
Definition replacement_rebuild.h:356
std::uint64_t bound
Definition replacement_rebuild.h:355
std::uint64_t quantum
Definition replacement_rebuild.h:355
bool tiny
Definition replacement_rebuild.h:358
std::uint64_t scan_price
Definition replacement_rebuild.h:355
bool building
Definition replacement_rebuild.h:357
std::deque< mutation > queue
Definition replacement_rebuild.h:352
rebuild(typed_world_type value)
Definition replacement_rebuild.h:359
std::uint64_t setup
Definition replacement_rebuild.h:355
Definition replacement_rebuild.h:129
static constexpr std::uint64_t tiny_record_limit
Definition replacement_rebuild.h:151
static std::uint64_t add(std::uint64_t a, std::uint64_t b)
Definition replacement_rebuild.h:390
void start(bool small)
Definition replacement_rebuild.h:459
static arrow_type clean(key_type const &key, state_type const &state)
Definition replacement_rebuild.h:394
static contribution_type put(key_type const &key, state_type const &value)
Definition replacement_rebuild.h:229
static std::uint64_t action_bound(std::uint64_t height)
Definition replacement_rebuild.h:452
std::uint64_t price() const
Definition replacement_rebuild.h:500
auto storage() const
Definition replacement_rebuild.h:197
world_type contribute(contribution_type input)
Definition replacement_rebuild.h:263
void candidate_work(auto &&operation)
Definition replacement_rebuild.h:507
static replacement_rebuild_engine from_snapshot(world_type source)
Definition replacement_rebuild.h:162
static void require(bool value, char const *message)
Definition replacement_rebuild.h:393
static replacement_rebuild_engine from_clean(typed_world_type source)
Definition replacement_rebuild.h:156
void writable() const
Definition replacement_rebuild.h:388
replacement_rebuild_work work_
Definition replacement_rebuild.h:365
void foreground_advance(std::uint64_t budget)
Definition replacement_rebuild.h:402
typed_detail::arrow_t< sort_type > arrow_type
Definition replacement_rebuild.h:136
static contribution_type command(mutation const &entry)
Definition replacement_rebuild.h:343
typed_detail::default_sort_t< P > sort_type
Definition replacement_rebuild.h:131
static std::uint64_t mul(std::uint64_t a, std::uint64_t b)
Definition replacement_rebuild.h:391
replacement_rebuild_engine(world_type value, Storage storage)
Definition replacement_rebuild.h:374
typed_detail::key_t< sort_type > key_type
Definition replacement_rebuild.h:134
static contribution_type change(key_type const &key, arrow_type const &arrow)
Definition replacement_rebuild.h:236
void active() const
Definition replacement_rebuild.h:387
std::unique_ptr< engine_type > foreground_
Definition replacement_rebuild.h:361
static constexpr std::uint64_t small_limit
Definition replacement_rebuild.h:150
typename engine_type::world_type typed_world_type
Definition replacement_rebuild.h:145
static auto batch()
Definition replacement_rebuild.h:227
auto convert_tiny(Source const &source)
Definition replacement_rebuild.h:526
bool pending() const noexcept
Definition replacement_rebuild.h:214
replacement_rebuild_work work() const noexcept
Definition replacement_rebuild.h:216
replacement_rebuild_engine(replacement_rebuild_engine const &)=delete
world_type published_
Definition replacement_rebuild.h:362
replacement_rebuild_status status() const
Definition replacement_rebuild.h:217
replacement_rebuild_engine & operator=(replacement_rebuild_engine const &)=delete
typed_detail::state_t< sort_type > state_type
Definition replacement_rebuild.h:135
static std::uint64_t reservation_work(std::uint64_t records)
Definition replacement_rebuild.h:241
std::uint64_t mutations_
Definition replacement_rebuild.h:364
replacement_rebuild_engine(world_type value)
Definition replacement_rebuild.h:368
bool admission_ready() const noexcept
Definition replacement_rebuild.h:215
void initialize_restore()
Definition replacement_rebuild.h:380
bool recovering_
Definition replacement_rebuild.h:366
typename engine_type::contribution_type contribution_type
Definition replacement_rebuild.h:147
void tiny_rebuild()
Definition replacement_rebuild.h:598
P policy_type
Definition replacement_rebuild.h:130
Family runtime_family
Definition replacement_rebuild.h:143
typename engine_type::runtime_type runtime_type
Definition replacement_rebuild.h:140
void grant(std::uint64_t amount)
Definition replacement_rebuild.h:496
std::unique_ptr< rebuild > job_
Definition replacement_rebuild.h:363
bool failed() const noexcept
Definition replacement_rebuild.h:186
void poison() noexcept
Definition replacement_rebuild.h:192
std::uint64_t base_
Definition replacement_rebuild.h:364
static replacement_rebuild_engine from_clean(typed_world_type source, Storage storage)
Definition replacement_rebuild.h:166
static std::uint64_t mass(typed_world_type const &state)
Definition replacement_rebuild.h:389
world_type publication() const
Definition replacement_rebuild.h:386
std::optional< world_type > advance(std::uint64_t budget)
Definition replacement_rebuild.h:323
replacement_rebuild_engine()
Definition replacement_rebuild.h:153
static std::uint64_t ceil(std::uint64_t a, std::uint64_t b)
Definition replacement_rebuild.h:392
static replacement_rebuild_engine from_snapshot(world_type source, Storage storage)
Definition replacement_rebuild.h:175
replacement_rebuild_engine(std::string schema)
Definition replacement_rebuild.h:154
static contribution_type erase(key_type const &key)
Definition replacement_rebuild.h:233
static session_reservation reservation(contribution_type const &input)
Definition replacement_rebuild.h:254
void apply(mutation entry)
Definition replacement_rebuild.h:408
static std::uint64_t tiny_conversion_bound()
Definition replacement_rebuild.h:522
world_type snapshot() const
Definition replacement_rebuild.h:185
bool failed_
Definition replacement_rebuild.h:366
replacement_rebuild_engine(replacement_rebuild_engine &&) noexcept=default
static constexpr bool charged_service
Definition replacement_rebuild.h:141
void rebase(world_type state)
Definition replacement_rebuild.h:203
replacement_world< P, A, Family > world_type
Definition replacement_rebuild.h:146
Definition replacement_rebuild.h:114
bool rebuilding
Definition replacement_rebuild.h:120
std::uint64_t horizon
Definition replacement_rebuild.h:116
bool scanning
Definition replacement_rebuild.h:120
std::uint64_t clean_rows
Definition replacement_rebuild.h:119
std::uint64_t replayed
Definition replacement_rebuild.h:116
std::uint64_t action_bound
Definition replacement_rebuild.h:117
std::uint64_t mutations
Definition replacement_rebuild.h:115
std::uint64_t initial_bound
Definition replacement_rebuild.h:117
std::uint64_t committed
Definition replacement_rebuild.h:118
std::uint64_t quantum
Definition replacement_rebuild.h:117
std::uint64_t queued
Definition replacement_rebuild.h:118
std::uint64_t admitted
Definition replacement_rebuild.h:116
std::uint64_t frozen_live
Definition replacement_rebuild.h:116
std::uint64_t credit
Definition replacement_rebuild.h:118
std::uint64_t source_records
Definition replacement_rebuild.h:119
std::uint64_t clean_base
Definition replacement_rebuild.h:115
Definition replacement_rebuild.h:106
std::uint64_t tiny_conversion_charged
Definition replacement_rebuild.h:112
std::uint64_t replayed
Definition replacement_rebuild.h:110
std::uint64_t reserved
Definition replacement_rebuild.h:109
std::uint64_t foreground_charged
Definition replacement_rebuild.h:108
std::uint64_t maximum_replay
Definition replacement_rebuild.h:111
std::uint64_t generations
Definition replacement_rebuild.h:107
std::uint64_t granted
Definition replacement_rebuild.h:109
std::uint64_t committed
Definition replacement_rebuild.h:109
std::uint64_t candidate_charged
Definition replacement_rebuild.h:108
std::uint64_t clean_rows
Definition replacement_rebuild.h:110
std::uint64_t tiny_indexes
Definition replacement_rebuild.h:112
std::uint64_t scan_records
Definition replacement_rebuild.h:110
std::uint64_t mutations
Definition replacement_rebuild.h:107
std::uint64_t maximum_handoff_mutations
Definition replacement_rebuild.h:111
std::uint64_t tiny_generations
Definition replacement_rebuild.h:112
Definition replacement_rebuild.h:88
std::shared_ptr< metadata_type const > metadata_
Definition replacement_rebuild.h:103
replacement_world(base_type value, std::uint64_t b=0, std::uint64_t u=0, bool active=false)
Definition replacement_rebuild.h:92
metadata_type const & metadata() const &noexcept
Definition replacement_rebuild.h:96
typename Family::snapshot_type runtime_snapshot
Definition replacement_rebuild.h:91
static replacement_world restore(runtime_snapshot data, metadata_type metadata, std::string_view schema)
Definition replacement_rebuild.h:97
Definition session.h:33
Definition typed_world.h:31
Definition typed_world.h:406
static typed_batch< P, A, Family > batch()
Definition typed_world.h:462
static typed_engine from_snapshot(world_type state)
Definition typed_world.h:433
static session_reservation reservation(contribution_type const &input)
Definition typed_world.h:477
static constexpr bool charged_service
Definition typed_world.h:417
typename Family::template runtime_type< compose_type > runtime_type
Definition typed_world.h:416
static std::uint64_t reservation_work(std::uint64_t records)
Definition typed_world.h:474
typed_world< P, A, Family > world_type
Definition typed_world.h:408
typed_contribution< P, A, Family > contribution_type
Definition typed_world.h:409
Definition typed_scan.h:33
bool has_row() const noexcept
Definition typed_scan.h:99
bool done() const noexcept
Definition typed_scan.h:98
std::uint64_t step(std::uint64_t budget)
Definition typed_scan.h:112
row_type take_row()
Definition typed_scan.h:102
std::uint64_t consumed() const noexcept
Definition typed_scan.h:101
Definition typed_world.h:185
std::string schema_id
Definition typed_world.h:188
std::uint64_t live_count
Definition typed_world.h:187
static typed_world_metadata decode(std::span< std::byte const > data)
Definition typed_world.h:202
std::vector< std::byte > encode() const
Definition typed_world.h:192
Definition typed_world.h:222
static typed_world restore(runtime_snapshot data, metadata_type metadata, std::string_view expected_schema)
Definition typed_world.h:238
runtime_snapshot const & runtime() const &noexcept
Definition typed_world.h:230
Scans one sort in an immutable typed world with chronological resolution.