Everett
Loading...
Searching...
No Matches
native_merge.h
Go to the documentation of this file.
1
13#pragma once
14
16
17#include <algorithm>
18#include <bit>
19#include <cstdint>
20#include <functional>
21#include <memory>
22#include <optional>
23#include <stdexcept>
24#include <type_traits>
25#include <utility>
26#include <vector>
27
28namespace everett {
30 bit_view operator()(bit_view, bit_view, bit_view newer) const { return newer; }
31 bit_view operator()(bit_view, bit_view newer) const { return newer; }
32 };
33
34 namespace native_merge_detail {
35 template <class T> T & composer(T & value) { return value; }
36 template <class T> T & composer(std::reference_wrapper<T> value) { return value.get(); }
37 template <class Compose, class Key> bool is_tombstone(Compose & compose, bit_view value, Key && key) {
38 auto & semantic = composer(compose);
39 if constexpr (requires { semantic.is_tombstone(value); })
40 return semantic.is_tombstone(value);
41 else if constexpr (requires { semantic.is_tombstone(bit_view{}, value); }) {
42 auto full = std::invoke(std::forward<Key>(key));
43 if constexpr (std::is_same_v<decltype(full), bit_view>) return semantic.is_tombstone(full, value);
44 else return semantic.is_tombstone(full.view(), value);
45 } else return false;
46 }
47 // Keep the predecessor as immutable literal spans, without copying its
48 // inherited prefix. A new retained boundary can precede the immediately
49 // preceding literal, so checking that one literal alone is insufficient.
50 // Each record adds at most one span; truncation removes spans permanently.
51 // Deep proper-prefix chains can retain one span per key unit.
52 template <class P> struct encoded_source {
54 : data_(view.bytes(), view.metadata().extent << P::unit_shift), cursor_(view) {
55 if (!done()) retain(cursor_.peek());
56 }
57 bool done() const noexcept { return cursor_.done(); }
58 profile_encoded_record const & peek() const & { return cursor_.peek(); }
60 std::uint64_t retained_bits() const { return peek().retained << P::unit_shift; }
62 bit_string result;
63 std::uint64_t begin = 0;
64 for (auto const & part : spans_) {
66 (part.end_units - begin) << P::unit_shift));
67 begin = part.end_units;
68 }
69 return result;
70 }
71
72 std::optional<bit_comparison> advance_comparison() {
74 if (done()) return std::nullopt;
75 auto const & next = cursor_.peek();
76 auto result = compare_successor(next.retained, next.suffix);
77 if (result.order >= 0)
78 throw std::invalid_argument("native merge inputs must have unique sorted keys");
79 retain(next);
80 return result;
81 }
82
83 private:
84 struct span {
85 std::uint64_t source_bit_offset;
86 std::uint64_t end_units;
87 };
88 static_assert(sizeof(span) == 16);
89 bit_comparison compare_successor(std::uint64_t retained, bit_view suffix) const {
90 auto first = spans_.size();
91 // Every crossed span beyond this boundary will be removed by retain;
92 // walking back is amortized over the one span added per input record.
93 while (first && spans_[first - 1].end_units > retained) --first;
94 auto begin = first ? spans_[first - 1].end_units : 0;
95 // Ordinary FC differs in the first remaining policy unit. Handle it
96 // directly; redundant controls fall through to the full fragment walk.
97 if (first != spans_.size() && !suffix.empty()) {
98 auto offset = spans_[first].source_bit_offset + ((retained - begin) << P::unit_shift);
99 auto before = profile_detail::load_bits(data_, offset, P::bits_per_unit);
100 auto after = profile_detail::load_bits(suffix, 0, P::bits_per_unit);
101 if (before != after) {
102 auto common = unsigned(std::countl_zero(before ^ after)) - (64 - P::bits_per_unit);
103 return {(retained << P::unit_shift) + common, before < after ? -1 : 1};
104 }
105 }
106 auto position = retained;
107 std::uint64_t compared = 0;
108 while (first != spans_.size() && compared != suffix.size()) {
109 auto const & fragment = spans_[first];
110 auto offset = fragment.source_bit_offset + ((position - begin) << P::unit_shift);
111 auto count = std::min((fragment.end_units - position) << P::unit_shift,
112 suffix.size() - compared);
113 auto result = compare_common_bits<typename P::architecture>(data_.subview(offset, count), suffix.subview(compared, count));
114 if (result.order) {
115 result.common_bits += (retained << P::unit_shift) + compared;
116 return result;
117 }
118 compared += count;
119 position = fragment.end_units;
120 begin = position;
121 ++first;
122 }
123 auto previous_units = spans_.empty() ? 0 : spans_.back().end_units;
124 auto previous_remaining = (previous_units - retained) << P::unit_shift;
125 return {(retained << P::unit_shift) + compared,
126 previous_remaining < suffix.size() ? -1 : previous_remaining > suffix.size() ? 1 : 0};
127 }
128 void retain(profile_encoded_record const & record) {
129 auto retained = record.retained;
130 while (!spans_.empty()) {
131 auto begin = spans_.size() > 1 ? spans_[spans_.size() - 2].end_units : 0;
132 if (begin >= retained) spans_.pop_back();
133 else {
134 if (spans_.back().end_units > retained) spans_.back().end_units = retained;
135 break;
136 }
137 }
138 // Every frame uses this one immutable backing. A span's logical start
139 // is the preceding endpoint, so truncation needs no copied length/view.
140 if (!record.suffix.empty()) spans_.push_back({record.suffix.offset(), record.key_units});
141 }
144 std::vector<span> spans_;
145 };
146 }
147
149 std::uint64_t keys = 0;
150 std::uint64_t input_records = 0;
151 };
152
153 // Merge two chronologically ordered native streams. Equal keys call
154 // compose(key, older_value, newer_value), or compose(older_value, newer_value)
155 // for value-only policies. Default composition is replacement. Key-aware
156 // policies retain materialized cursors; the other cases forward encoded
157 // literals and validate strict input order through immutable prefix spans.
158 // A custom policy accepting both forms uses the key-aware form.
159 // Each input must have unique sorted keys. Input owners stay pinned,
160 // including after a decoding/composition failure.
161 // Callback keys borrow cursor scratch. A policy retaining a value view must
162 // retain its source owner too; reassignment can release earlier source pins.
163 // Construction parses the first record from each nonempty source before any
164 // step budget is charged; key-aware policies also reconstruct those keys.
165 // One step unit handles one distinct key and at most two input records.
166 // Key bytes, composition work, output allocation and final EF work are extra.
167 // An alternate Output consumes append(retained,literal,value) synchronously,
168 // encodes policy P, and supplies size/common_value_width/finished/failed/finish.
169 // It starts empty and active; finished() and failed() are noexcept. The caller
170 // grants exclusive sink use to this builder, and keeps any referenced sink
171 // alive. A finish failure is retryable only when the sink reports !failed().
172 template <class P, class Native = profile_array<P>, class Compose = replace_native_value,
173 class Output = profile_detail::native_output<P>>
175 static_assert(std::is_same_v<typename Native::policy_type, P>);
176 static_assert(std::is_same_v<typename Output::policy_type, P>);
177 using policy_type = P;
178 using source_type = Native;
179 using output_type = Output;
180 static_assert(noexcept(std::declval<Output const &>().finished()));
181 static_assert(noexcept(std::declval<Output const &>().failed()));
182 using source_pointer = std::shared_ptr<Native const>;
183 static constexpr bool encoded_keys = std::is_same_v<Compose, replace_native_value> ||
184 (!std::is_invocable_v<Compose &, bit_view, bit_view, bit_view> &&
185 std::is_invocable_v<Compose &, bit_view, bit_view>);
186 static_assert(encoded_keys || std::is_invocable_v<Compose &, bit_view, bit_view, bit_view>,
187 "native composition accepts (key, older, newer) or (older, newer)");
188
189 native_merge_builder(source_pointer older, source_pointer newer, Compose compose = {},
190 std::optional<std::uint64_t> common_value_width = P::value_width)
191 requires std::is_constructible_v<Output, std::optional<std::uint64_t>>
192 : native_merge_builder(Output(common_value_width), std::move(older), std::move(newer), std::move(compose)) {}
193 native_merge_builder(Output output, source_pointer older, source_pointer newer, Compose compose = {})
194 : older_(checked(std::move(older))), newer_(checked(std::move(newer))),
195 older_cursor_(older_->view()), newer_cursor_(newer_->view()),
196 writer_(checked_output(std::move(output))), compose_(std::move(compose)) {}
201 requires (std::is_move_assignable_v<Compose> && std::is_move_assignable_v<Output>) {
202 if (this != &other) {
203 // Composition policies may throw during transfer. Keep a partially
204 // assigned destination poisoned until every state field agrees.
205 failed_ = true;
206 older_ = std::move(other.older_); newer_ = std::move(other.newer_);
207 older_cursor_ = std::move(other.older_cursor_); newer_cursor_ = std::move(other.newer_cursor_);
208 writer_ = std::move(other.writer_); compose_ = std::move(other.compose_);
209 progress_ = other.progress_;
210 older_prefix_ = other.older_prefix_; newer_prefix_ = other.newer_prefix_;
211 failed_ = other.failed_;
212 }
213 return *this;
214 }
215
216 bool done() const noexcept {
217 return older_ && newer_ && !failed_ && older_cursor_.done() && newer_cursor_.done();
218 }
219 bool failed() const noexcept { return failed_; }
220 bool finished() const noexcept { return writer_.finished(); }
221 native_merge_progress progress() const noexcept { return progress_; }
222 source_pointer older_source() const noexcept { return older_; }
223 source_pointer newer_source() const noexcept { return newer_; }
224
225 native_merge_progress step(std::uint64_t key_budget = 1) {
228 try {
229 while (work.keys < key_budget && !done()) {
230 auto comparison = compare_heads();
231 auto order = comparison.order;
232 auto consumed = order == 0 ? 2u : 1u;
233 auto total_keys = profile_detail::add(progress_.keys, 1);
235 if (order < 0) {
236 auto item = older_cursor_.peek();
237 append(older_cursor_, item, item.value, older_prefix_, older_cursor_.retained_bits());
238 newer_prefix_ = comparison.common_bits >> P::unit_shift;
240 } else if (order > 0) {
241 auto item = newer_cursor_.peek();
242 append(newer_cursor_, item, item.value, newer_prefix_, newer_cursor_.retained_bits());
243 older_prefix_ = comparison.common_bits >> P::unit_shift;
245 } else {
246 auto older = older_cursor_.peek(), newer = newer_cursor_.peek();
247 auto value = [&] {
248 if constexpr (encoded_keys) return std::invoke(compose_, older.value, newer.value);
249 else return std::invoke(compose_, older.key.prefix, older.value, newer.value);
250 }();
251 auto limit = std::min(older_cursor_.retained_bits(), newer_cursor_.retained_bits());
252 auto emit = [&](auto const & source, auto const & item) {
253 if constexpr (std::is_same_v<decltype(value), bit_view>)
254 append(source, item, value, older_prefix_, limit);
255 else append(source, item, value.view(), older_prefix_, limit);
256 };
257 if constexpr (encoded_keys) {
258 // Equal keys have interchangeable literal donors. Use the one
259 // that already carries the conservative tombstone boundary,
260 // avoiding reconstruction of an older, shorter literal.
261 if (newer.retained < older.retained) emit(newer_cursor_, newer);
262 else emit(older_cursor_, older);
263 } else emit(older_cursor_, older);
266 }
267 progress_ = {total_keys, total_records};
268 ++work.keys;
269 work.input_records += consumed;
270 }
271 } catch (...) {
272 failed_ = true;
273 throw;
274 }
275 return work;
276 }
277
278 auto finish() {
280 if (!done()) throw std::logic_error("native merge has unconsumed inputs");
281 try { return writer_.finish(); }
282 catch (...) {
283 // In-memory EF allocation failure is retryable. A sink that has
284 // partially written final output must report its irreversible failure.
285 failed_ = writer_.failed();
286 throw;
287 }
288 }
289
290 private:
291 using cursor_type = std::conditional_t<encoded_keys, native_merge_detail::encoded_source<P>,
293 static bit_view suffix(auto const & item, std::uint64_t retained) {
294 if constexpr (encoded_keys) {
295 if (retained < item.retained)
296 throw std::invalid_argument("native merge output prefix precedes input prefix");
297 auto first = (retained - item.retained) << P::unit_shift;
298 return item.suffix.subview(first, item.suffix.size() - first);
299 } else {
300 auto first = retained << P::unit_shift;
301 return item.key.prefix.subview(first, item.key.prefix.size() - first);
302 }
303 }
304 void append(cursor_type const & source, auto const & item, bit_view value,
305 std::uint64_t retained, std::uint64_t limit_bits) {
307 if constexpr (encoded_keys) return source.materialize();
308 else return item.key.prefix;
309 })) retained = std::min(retained, limit_bits >> P::unit_shift);
310 if constexpr (encoded_keys) {
311 if (retained < item.retained) {
312 auto key = source.materialize();
313 auto first = retained << P::unit_shift;
314 writer_.append(retained, key.view().subview(first, key.bit_size - first), value);
315 return;
316 }
317 }
318 writer_.append(retained, suffix(item, retained), value);
319 }
320 // Both heads follow the last emitted key p. The head sharing more of p
321 // sorts first. Equal LCPs need only a suffix comparison from that boundary.
323 if (older_cursor_.done()) return {0, 1};
324 if (newer_cursor_.done()) return {0, -1};
326 return {std::min(older_prefix_, newer_prefix_) << P::unit_shift,
327 older_prefix_ > newer_prefix_ ? -1 : 1};
328 auto start = older_prefix_ << P::unit_shift;
329 auto result = compare_common_bits<typename P::architecture>(suffix(older_cursor_.peek(), older_prefix_),
331 result.common_bits += start;
332 return result;
333 }
334 static std::uint64_t advance(cursor_type & cursor) {
335 auto comparison = cursor.advance_comparison();
336 if (!comparison) return 0;
337 if (comparison->order >= 0)
338 throw std::invalid_argument("native merge inputs must have unique sorted keys");
339 return comparison->common_bits >> P::unit_shift;
340 }
341 static Output checked_output(Output output) {
342 if (output.size() || output.finished() || output.failed())
343 throw std::invalid_argument("native merge output must be empty and active");
344 return output;
345 }
347 if (!source) throw std::invalid_argument("null native merge input");
348 return source;
349 }
350 void require_active() const {
351 if (!older_ || !newer_ || failed_ || writer_.finished())
352 throw std::logic_error("native merge is inactive");
353 }
358 Output writer_;
359 Compose compose_;
361 std::uint64_t older_prefix_ = 0;
362 std::uint64_t newer_prefix_ = 0;
363 bool failed_ = false;
364 };
365}
bool is_tombstone(Compose &compose, bit_view value, Key &&key)
Definition native_merge.h:37
T & composer(T &value)
Definition native_merge.h:35
std::uint64_t add(std::uint64_t a, std::uint64_t b)
Definition profile.h:39
std::uint64_t load_bits(bit_view data, std::uint64_t first, unsigned width) noexcept
Definition profile.h:93
void append(bit_string &target, bit_view source)
Definition profile.h:480
Definition active_engine.h:18
Builds native front-coded profiles incrementally with explicit value widths.
Definition profile.h:209
Definition profile.h:166
bit_view view() const &
Definition profile.h:178
std::uint64_t bit_size
Definition profile.h:168
Definition profile.h:56
std::uint64_t size() const noexcept
Definition profile.h:64
std::uint64_t offset() const noexcept
Definition profile.h:66
bool empty() const noexcept
Definition profile.h:65
bit_view subview(std::uint64_t first, std::uint64_t count) const
Definition profile.h:73
Definition native_merge.h:174
source_pointer older_
Definition native_merge.h:354
std::uint64_t older_prefix_
Definition native_merge.h:361
Native source_type
Definition native_merge.h:178
std::uint64_t newer_prefix_
Definition native_merge.h:362
native_merge_builder(Output output, source_pointer older, source_pointer newer, Compose compose={})
Definition native_merge.h:193
native_merge_progress step(std::uint64_t key_budget=1)
Definition native_merge.h:225
std::shared_ptr< Native const > source_pointer
Definition native_merge.h:182
Output output_type
Definition native_merge.h:179
cursor_type older_cursor_
Definition native_merge.h:356
native_merge_builder(native_merge_builder const &)=delete
auto finish()
Definition native_merge.h:278
P policy_type
Definition native_merge.h:177
std::conditional_t< encoded_keys, native_merge_detail::encoded_source< P >, profile_cursor< P, stream_role::native > > cursor_type
Definition native_merge.h:292
native_merge_builder & operator=(native_merge_builder &&other)
Definition native_merge.h:200
static source_pointer checked(source_pointer source)
Definition native_merge.h:346
Output writer_
Definition native_merge.h:358
bool done() const noexcept
Definition native_merge.h:216
bool finished() const noexcept
Definition native_merge.h:220
native_merge_builder(source_pointer older, source_pointer newer, Compose compose={}, std::optional< std::uint64_t > common_value_width=P::value_width)
Definition native_merge.h:189
native_merge_progress progress() const noexcept
Definition native_merge.h:221
bool failed() const noexcept
Definition native_merge.h:219
void require_active() const
Definition native_merge.h:350
source_pointer newer_source() const noexcept
Definition native_merge.h:223
void append(cursor_type const &source, auto const &item, bit_view value, std::uint64_t retained, std::uint64_t limit_bits)
Definition native_merge.h:304
source_pointer newer_
Definition native_merge.h:355
native_merge_progress progress_
Definition native_merge.h:360
Compose compose_
Definition native_merge.h:359
static constexpr bool encoded_keys
Definition native_merge.h:183
static bit_view suffix(auto const &item, std::uint64_t retained)
Definition native_merge.h:293
bit_comparison compare_heads() const
Definition native_merge.h:322
source_pointer older_source() const noexcept
Definition native_merge.h:222
static std::uint64_t advance(cursor_type &cursor)
Definition native_merge.h:334
native_merge_builder(native_merge_builder &&)=default
bool failed_
Definition native_merge.h:363
native_merge_builder & operator=(native_merge_builder const &)=delete
static Output checked_output(Output output)
Definition native_merge.h:341
cursor_type newer_cursor_
Definition native_merge.h:357
std::uint64_t source_bit_offset
Definition native_merge.h:85
std::uint64_t end_units
Definition native_merge.h:86
encoded_source(profile_view< P > view)
Definition native_merge.h:53
bit_comparison compare_successor(std::uint64_t retained, bit_view suffix) const
Definition native_merge.h:89
profile_encoded_cursor< P > cursor_
Definition native_merge.h:143
void retain(profile_encoded_record const &record)
Definition native_merge.h:128
bit_string materialize() const
Definition native_merge.h:61
profile_encoded_record const & peek() const &&=delete
std::uint64_t retained_bits() const
Definition native_merge.h:60
std::optional< bit_comparison > advance_comparison()
Definition native_merge.h:72
std::vector< span > spans_
Definition native_merge.h:144
bool done() const noexcept
Definition native_merge.h:57
bit_view data_
Definition native_merge.h:142
profile_encoded_record const & peek() const &
Definition native_merge.h:58
Definition native_merge.h:148
std::uint64_t input_records
Definition native_merge.h:150
std::uint64_t keys
Definition native_merge.h:149
Definition profile.h:1066
bool done() const noexcept
Definition profile.h:1023
void advance()
Definition profile.h:1031
profile_encoded_record const & peek() const &
Definition profile.h:1025
Definition profile.h:308
bit_view suffix
Definition profile.h:312
std::uint64_t retained
Definition profile.h:309
std::uint64_t key_units
Definition profile.h:310
Definition native_merge.h:29
bit_view operator()(bit_view, bit_view, bit_view newer) const
Definition native_merge.h:30
bit_view operator()(bit_view, bit_view newer) const
Definition native_merge.h:31