Everett
Loading...
Searching...
No Matches
sort_profile_adaptive.h
Go to the documentation of this file.
1
12#pragma once
13
17
18#include <variant>
19
20namespace everett {
21 namespace sort_profile_detail {
22 // A factory reserves its identity before constructing the private stream.
23 // Until the first flush, only the existing fixed bit-sink buffer is used.
24 template <class P, class Ops, class Factory> struct lazy_native_stream {
26 explicit lazy_native_stream(Factory factory) : factory_(std::move(factory)) {}
27 bool spilled() const noexcept { return bool(stream_); }
28 bool failed() const noexcept { return failed_ || (stream_ && stream_->failed()); }
29 std::uint64_t body_bytes() const noexcept { return stream_ ? stream_->body_bytes() : 192; }
30 void start() {
31 if (failed()) throw std::logic_error("failed adaptive native stream");
32 if (stream_) return;
33 try {
34 stream_ = factory_();
35 if (!stream_) throw std::logic_error("null adaptive native stream");
36 } catch (...) { poison(); throw; }
37 }
38 void append(std::span<std::byte const> bytes) {
39 if (bytes.empty()) return;
40 try { start(); stream_->append(bytes); }
41 catch (...) { poison(); throw; }
42 }
43 object_seal_receipt finish(file_header<P> const & header, std::span<std::byte const> directory) {
44 try { start(); return stream_->finish(header, directory); }
45 catch (...) { poison(); throw; }
46 }
47 object_write_paths const & paths() const & {
48 if (!stream_) throw std::logic_error("native output has not spilled");
49 return stream_->paths();
50 }
51 object_write_paths const & paths() const && = delete;
52 void poison() noexcept {
53 failed_ = true;
54 if constexpr (requires { { factory_.poison() } noexcept; }) factory_.poison();
55 }
56 private:
57 Factory factory_;
58 std::unique_ptr<stream_type> stream_;
59 bool failed_ = false;
60 };
61
62 // Compose may expand values, so it cannot bound payload bytes in advance.
63 // It cannot introduce additional records or selector codes. These cheap
64 // source-metadata bounds reject a clearly oversized navigation structure;
65 // actual finished allocation capacities are checked again before retention.
66 template <class P, class View> bool small_native_metadata(View const & a, View const & b, std::size_t limit) {
67 auto use = [&](std::uint64_t count, std::size_t width) {
68 if (count > limit / width) return false;
69 limit -= std::size_t(count) * width; return true;
70 };
71 if (b.size() > std::numeric_limits<std::uint64_t>::max() - a.size()) return false;
72 auto records = a.size() + b.size();
73 auto blocks = records / P::codec_block_size + (records % P::codec_block_size != 0);
74 // EF low/high/select storage and packed selector seeds. This intentionally
75 // allows slack; it is a preflight filter, not allocator measurement.
76 return use(1, 512) && use(blocks + 1, 64) &&
77 use(a.dictionary_size() + 1, 16) && use(b.dictionary_size() + 1, 16) &&
78 use(profile_detail::byte_count(a.dictionary().size()), 2) &&
79 use(profile_detail::byte_count(b.dictionary().size()), 2);
80 }
81 }
82
83 template <class P, class Selector, class Ops, class Factory> struct sort_profile_adaptive_writer {
85 using result_type = std::variant<std::shared_ptr<array_type const>, object_seal_receipt>;
88 sort_profile_adaptive_writer(Factory factory, output_budget budget, std::size_t ceiling, bool metadata_fits)
89 : stream_(std::move(factory)), sink_(stream_), budget_(std::move(budget)), ceiling_(ceiling) {
90 if (!ceiling_ || !budget_.limit() || !metadata_fits) stream_.start();
91 }
94 bool failed() const noexcept { return failed_ || stream_.failed(); }
95 bool finished() const noexcept { return finished_; }
96 bool spilled() const noexcept { return stream_.spilled(); }
97 std::uint64_t size() const noexcept { return encoder_.size(); }
98 object_write_paths const & paths() const & { return stream_.paths(); }
99 object_write_paths const & paths() const && = delete;
100 void append_frame(sort_profile_frame const & frame, std::span<bit_view const> key,
101 std::uint64_t common, bit_view value) {
103 try { encoder_.append_frame(sink_, frame, key, common, value); }
104 catch (...) { failed_ = true; stream_.poison(); throw; }
105 }
108 try {
109 encoder_.finish(sink_.position());
110 if (!stream_.spilled()) {
112 auto const & ef = encoder_.offsets;
113 auto bytes = sizeof(array_type) + sizeof(output_budget::lease) + data.bytes.capacity() + encoder_.dictionary.bytes.capacity() +
114 encoder_.seeds.bytes.capacity() + 8 * (encoder_.dictionary_offsets.capacity() +
115 ef.low.capacity() + ef.high.capacity() + ef.sparse.capacity()) +
116 sizeof(elias_fano_sample) * ef.samples.capacity();
117 if (bytes <= ceiling_) if (auto lease = budget_.try_acquire(bytes)) {
118 auto result = output_budget::attach(encoder_.take(std::move(data)), std::move(*lease));
119 finished_ = true; return result;
120 }
121 }
122 auto result = sort_profile_detail::seal<P>(encoder_, sink_, stream_);
123 finished_ = true; return result;
124 } catch (...) { failed_ = true; stream_.poison(); throw; }
125 }
126 private:
130 std::size_t ceiling_;
132 bool failed_ = false, finished_ = false;
133 void require_active() const { if (failed() || finished_) throw std::logic_error("inactive adaptive native writer"); }
134 };
135
136 template <class P, class Native, class Compose, class Selector, class Ops, class Factory>
138 using source_pointer = std::shared_ptr<Native const>;
140 struct output_ref {
142 bool failed() const noexcept { return output->failed(); }
143 void append_frame(sort_profile_frame const & frame, std::span<bit_view const> key,
144 std::uint64_t common, bit_view value) { output->append_frame(frame, key, common, value); }
145 auto finish() { return output->finish(); }
146 };
148 sort_profile_adaptive_merge(Factory factory, output_budget budget, std::size_t ceiling,
149 source_pointer older, source_pointer newer, Compose compose)
150 : output_(std::move(factory), std::move(budget), ceiling, metadata_fits(older, newer, ceiling)),
151 builder_({&output_}, std::move(older), std::move(newer), std::move(compose)) {}
154 bool done() const noexcept { return builder_.done(); }
155 bool failed() const noexcept { return builder_.failed(); }
156 bool finished() const noexcept { return builder_.finished(); }
157 bool spilled() const noexcept { return output_.spilled(); }
158 native_merge_progress progress() const noexcept { return builder_.progress(); }
159 std::uint64_t materialized_keys() const noexcept { return builder_.materialized_keys(); }
160 native_merge_progress step(std::uint64_t budget = 1) { return builder_.step(budget); }
161 auto finish() { return builder_.finish(); }
162 object_write_paths const & paths() const & { return output_.paths(); }
163 object_write_paths const & paths() const && = delete;
164 private:
167 static bool metadata_fits(source_pointer const & a, source_pointer const & b, std::size_t ceiling) {
168 if (!a || !b) throw std::invalid_argument("null adaptive merge source");
169 return sort_profile_detail::small_native_metadata<P>(a->view(), b->view(), ceiling);
170 }
171 };
172}
std::uint64_t byte_count(std::uint64_t bits) noexcept
Definition profile.h:50
bool small_native_metadata(View const &a, View const &b, std::size_t limit)
Definition sort_profile_adaptive.h:66
Definition active_engine.h:18
Shares an allowance for encoded outputs retained in memory.
Streams sort-owned KV03 payloads with bounded buffering and shared framing.
Merges sort-owned records while retaining inherited keys as borrowed spans.
static bit_string copy(bit_view source)
Definition profile.h:181
Definition profile.h:56
Definition elias_fano.h:195
Definition file.h:34
Definition native_merge.h:148
Definition object_writer.h:88
Definition object_writer.h:53
Definition output_budget.h:34
Definition output_budget.h:26
std::size_t limit() const noexcept
Definition output_budget.h:69
std::optional< lease > try_acquire(std::size_t bytes) const noexcept
Definition output_budget.h:71
static std::shared_ptr< T const > attach(T value, lease allocation)
Definition output_budget.h:83
Definition sort_profile_adaptive.h:140
void append_frame(sort_profile_frame const &frame, std::span< bit_view const > key, std::uint64_t common, bit_view value)
Definition sort_profile_adaptive.h:143
output_type * output
Definition sort_profile_adaptive.h:141
bool failed() const noexcept
Definition sort_profile_adaptive.h:142
auto finish()
Definition sort_profile_adaptive.h:145
Definition sort_profile_adaptive.h:137
object_write_paths const & paths() const &&=delete
sort_profile_adaptive_merge(Factory factory, output_budget budget, std::size_t ceiling, source_pointer older, source_pointer newer, Compose compose)
Definition sort_profile_adaptive.h:148
static bool metadata_fits(source_pointer const &a, source_pointer const &b, std::size_t ceiling)
Definition sort_profile_adaptive.h:167
native_merge_progress step(std::uint64_t budget=1)
Definition sort_profile_adaptive.h:160
native_merge_progress progress() const noexcept
Definition sort_profile_adaptive.h:158
bool finished() const noexcept
Definition sort_profile_adaptive.h:156
std::shared_ptr< Native const > source_pointer
Definition sort_profile_adaptive.h:138
output_type output_
Definition sort_profile_adaptive.h:165
bool done() const noexcept
Definition sort_profile_adaptive.h:154
bool spilled() const noexcept
Definition sort_profile_adaptive.h:157
builder_type builder_
Definition sort_profile_adaptive.h:166
sort_profile_adaptive_merge(sort_profile_adaptive_merge const &)=delete
auto finish()
Definition sort_profile_adaptive.h:161
sort_profile_adaptive_merge & operator=(sort_profile_adaptive_merge const &)=delete
std::uint64_t materialized_keys() const noexcept
Definition sort_profile_adaptive.h:159
bool failed() const noexcept
Definition sort_profile_adaptive.h:155
object_write_paths const & paths() const &
Definition sort_profile_adaptive.h:162
Definition sort_profile_adaptive.h:83
bool finished() const noexcept
Definition sort_profile_adaptive.h:95
sort_profile_array< P, Selector > array_type
Definition sort_profile_adaptive.h:84
sort_profile_detail::encoder< P, Selector > encoder_
Definition sort_profile_adaptive.h:131
sort_profile_adaptive_writer & operator=(sort_profile_adaptive_writer const &)=delete
sort_profile_adaptive_writer(Factory factory, output_budget budget, std::size_t ceiling, bool metadata_fits)
Definition sort_profile_adaptive.h:88
bool spilled() const noexcept
Definition sort_profile_adaptive.h:96
void require_active() const
Definition sort_profile_adaptive.h:133
std::size_t ceiling_
Definition sort_profile_adaptive.h:130
bool finished_
Definition sort_profile_adaptive.h:132
std::variant< std::shared_ptr< array_type const >, object_seal_receipt > result_type
Definition sort_profile_adaptive.h:85
output_budget budget_
Definition sort_profile_adaptive.h:129
object_write_paths const & paths() const &&=delete
object_write_paths const & paths() const &
Definition sort_profile_adaptive.h:98
bool failed_
Definition sort_profile_adaptive.h:132
bool failed() const noexcept
Definition sort_profile_adaptive.h:94
stream_type stream_
Definition sort_profile_adaptive.h:127
sink_type sink_
Definition sort_profile_adaptive.h:128
void append_frame(sort_profile_frame const &frame, std::span< bit_view const > key, std::uint64_t common, bit_view value)
Definition sort_profile_adaptive.h:100
result_type finish()
Definition sort_profile_adaptive.h:106
sort_profile_adaptive_writer(sort_profile_adaptive_writer const &)=delete
std::uint64_t size() const noexcept
Definition sort_profile_adaptive.h:97
Definition sort_profile.h:405
Definition sort_profile.h:443
std::uint64_t position() const noexcept
Definition sort_profile_file_writer.h:22
bit_view buffered_payload() const &noexcept
Definition sort_profile_file_writer.h:26
Definition sort_profile_adaptive.h:24
bool failed_
Definition sort_profile_adaptive.h:59
void append(std::span< std::byte const > bytes)
Definition sort_profile_adaptive.h:38
std::uint64_t body_bytes() const noexcept
Definition sort_profile_adaptive.h:29
bool failed() const noexcept
Definition sort_profile_adaptive.h:28
object_write_paths const & paths() const &
Definition sort_profile_adaptive.h:47
std::unique_ptr< stream_type > stream_
Definition sort_profile_adaptive.h:58
object_write_paths const & paths() const &&=delete
lazy_native_stream(Factory factory)
Definition sort_profile_adaptive.h:26
bool spilled() const noexcept
Definition sort_profile_adaptive.h:27
void poison() noexcept
Definition sort_profile_adaptive.h:52
void start()
Definition sort_profile_adaptive.h:30
object_seal_receipt finish(file_header< P > const &header, std::span< std::byte const > directory)
Definition sort_profile_adaptive.h:43
Factory factory_
Definition sort_profile_adaptive.h:57
Definition sort_profile.h:163
std::uint64_t materialized_keys() const noexcept
Definition sort_profile_merge.h:78
auto finish()
Definition sort_profile_merge.h:111
bool failed() const noexcept
Definition sort_profile_merge.h:75
bool done() const noexcept
Definition sort_profile_merge.h:74
native_merge_progress step(std::uint64_t budget=1)
Definition sort_profile_merge.h:79
bool finished() const noexcept
Definition sort_profile_merge.h:76
native_merge_progress progress() const noexcept
Definition sort_profile_merge.h:77