Everett
Loading...
Searching...
No Matches
cola_adaptive_index.h
Go to the documentation of this file.
1
12#pragma once
13
16#include <variant>
17
18namespace everett {
19 namespace cola_detail {
20 template <class P, class Ops, class SpoolOps> struct index_destination {
24 index_destination(std::filesystem::path root, object_id id, object_attempt_id attempt,
25 cola_file_dependencies deps, Ops & ops, SpoolOps & spool_ops)
26 : dependencies(std::move(deps)), stream(std::move(root), std::move(id), std::move(attempt),
27 file_kind::fractional_index, cola_section_detail::directory_bytes, ops),
28 spool(stream.paths().private_output.string() + ".secondary", spool_ops) {}
29 };
30
31 // Factory acquires a physical destination only at the first spill. It owns
32 // any borrowed Ops, so it outlives the destination and both staged sinks.
33 template <class P, class Native, class Main, class Ops, class SpoolOps, class Factory>
37 using result_type = std::variant<std::shared_ptr<index_type const>, object_seal_receipt>;
38 struct route_stream {
40 unsigned route;
41 void append(std::span<std::byte const> bytes) {
42 if (bytes.empty()) return;
43 auto & target = owner->destination();
44 if (!route) target.stream.append(bytes); else target.spool.append(bytes);
45 }
46 std::uint64_t body_bytes() const noexcept {
48 return !route ? owner->destination_->stream.body_bytes() : owner->destination_->spool.body_bytes();
49 }
50 };
52 adaptive_index_output(Factory factory, output_budget budget, std::size_t limit)
53 : factory_(std::move(factory)), budget_(std::move(budget)), limit_(limit),
54 routes_{{{this, 0}, {this, 1}}}, main_(routes_[0]), secondary_(routes_[1]) {
55 if (!limit_ || !budget_.limit()) (void)destination();
56 }
59 bool failed() const noexcept { return failed_ || (destination_ && destination_->stream.failed()); }
60 bool finished() const noexcept { return finished_; }
61 bool spilled() const noexcept { return bool(destination_); }
62 Factory const & factory() const noexcept { return factory_; }
63 object_write_paths const & paths() const & {
64 if (!destination_) throw std::logic_error("index has no physical output");
65 return destination_->stream.paths();
66 }
67 object_write_paths const & paths() const && = delete;
68 std::uint64_t spooled_bytes() const noexcept { return destination_ ? destination_->spool.body_bytes() : 0; }
69 void append_known(unsigned route, bit_view key, std::uint64_t common) {
70 require_active();
71 try {
72 if (!route) profiles_[0].append(main_, key, common);
73 else if (route == 1) profiles_[1].append(secondary_, key, common);
74 else throw std::out_of_range("adaptive index route");
75 } catch (...) { poison(); throw; }
76 }
79 require_active();
80 try {
81 profiles_[0].finish_metadata(main_.position());
82 profiles_[1].finish_metadata(secondary_.position());
83 if (!spilled()) {
84 // Copy only the bounded staged payload; all sparse arrays move.
85 std::array<bit_string, 2> payload{bit_string::copy(main_.buffered_payload()),
86 bit_string::copy(secondary_.buffered_payload())};
87 auto charge = retained_bytes(payload, metadata);
88 if (charge <= limit_) if (auto lease = budget_.try_acquire(charge)) {
90 std::array<typename index_type::borrowed_array, 2> borrowed{
91 adoption::adopt(std::move(payload[0].bytes), std::move(profiles_[0].offsets), profiles_[0].metadata),
92 adoption::adopt(std::move(payload[1].bytes), std::move(profiles_[1].offsets), profiles_[1].metadata)};
93 auto built = index_output<P, Native, Main>::adopt(std::move(native), std::move(main), std::move(secondary),
94 std::move(borrowed), std::move(metadata));
95 auto owned = output_budget::attach(std::move(built), std::move(*lease));
96 finished_ = true; return owned;
97 }
98 }
99 auto & target = destination();
100 if (!native || bool(main) != bool(target.dependencies.main) || bool(secondary) != bool(target.dependencies.secondary))
101 throw std::invalid_argument("adaptive COLA dependency shape");
102 main_.finish_payload(); secondary_.finish_payload();
103 auto receipt = seal_file_index<P>(target.dependencies, profiles_, metadata, native->size(),
104 main_, target.stream, target.spool);
105 finished_ = true; return receipt;
106 } catch (...) { poison(); throw; }
107 }
108 private:
109 Factory factory_;
111 std::size_t limit_;
112 std::unique_ptr<destination_type> destination_;
113 std::array<route_stream, 2> routes_;
114 sink_type main_, secondary_;
115 std::array<borrowed_file_state<P>, 2> profiles_;
116 bool failed_ = false, finished_ = false;
117 void poison() noexcept {
118 failed_ = true;
119 if constexpr (requires { { factory_.poison() } noexcept; }) factory_.poison();
120 }
122 if (!destination_) {
123 destination_ = factory_();
124 if (!destination_) throw std::logic_error("null adaptive index destination");
125 }
126 return *destination_;
127 }
128 void require_active() const { if (failed() || finished()) throw std::logic_error("inactive adaptive index output"); }
129 std::size_t retained_bytes(std::array<bit_string, 2> const & payload, index_metadata<P> const & metadata) const {
130 std::uint64_t bytes = sizeof(index_type) + sizeof(output_budget::lease);
131 auto add = [&](auto const & values) {
132 bytes = profile_detail::add(bytes, profile_detail::multiply(values.capacity(), sizeof(typename std::remove_cvref_t<decltype(values)>::value_type)));
133 };
134 for (unsigned route = 0; route != 2; ++route) {
135 add(payload[route].bytes);
136 auto const & offsets = profiles_[route].offsets;
137 add(offsets.low); add(offsets.high); add(offsets.samples); add(offsets.sparse);
138 add(metadata.ranks[route].classes); add(metadata.ranks[route].checkpoints);
139 add(metadata.flags[route]); add(metadata.cuts[route]);
140 }
141 if (bytes > std::numeric_limits<std::size_t>::max()) throw std::length_error("retained index extent");
142 return static_cast<std::size_t>(bytes);
143 }
144 };
145 }
146
147 template <class P, class Native, class Main, class Factory,
148 class Ops = posix_object_ops, class SpoolOps = posix_index_spool_ops>
154 cola_adaptive_index_builder(Factory factory, output_budget budget, std::size_t limit,
155 native_pointer native, main_pointer main = {}, native_pointer secondary = {})
156 : output_(std::move(factory), std::move(budget), limit),
157 builder_({&output_}, std::move(native), std::move(main), std::move(secondary)) {}
160 bool done() const noexcept { return builder_.done(); }
161 bool failed() const noexcept { return builder_.failed(); }
162 bool finished() const noexcept { return builder_.finished(); }
163 bool spilled() const noexcept { return output_.spilled(); }
164 std::uint64_t size() const noexcept { return builder_.size(); }
165 native_pointer native_owner() const noexcept { return builder_.native_owner(); }
166 main_pointer main_target() const noexcept { return builder_.main_target(); }
167 native_pointer secondary_target() const noexcept { return builder_.secondary_target(); }
168 std::uint64_t step(std::uint64_t budget) { return builder_.step(budget); }
169 auto finish() { return builder_.finish(); }
170 Factory const & factory() const noexcept { return output_.factory(); }
171 std::uint64_t spooled_bytes() const noexcept { return output_.spooled_bytes(); }
172 object_write_paths const & paths() const & { return output_.paths(); }
173 object_write_paths const & paths() const && = delete;
174 private:
175 output_type output_;
176 builder_type builder_;
177 };
178}
Streams dual-target IX03 indexes with one private secondary payload spool.
unsigned route(unsigned value)
Definition cola_index.h:37
constexpr std::size_t directory_bytes
Definition cola_sections.h:31
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
Definition active_engine.h:18
file_kind
Definition object_path.h:25
Shares an allowance for encoded outputs retained in memory.
static bit_string copy(bit_view source)
Definition profile.h:181
Definition profile.h:56
Definition cola_adaptive_index.h:149
bool done() const noexcept
Definition cola_adaptive_index.h:160
typename builder_type::main_pointer main_pointer
Definition cola_adaptive_index.h:153
native_pointer native_owner() const noexcept
Definition cola_adaptive_index.h:165
object_write_paths const & paths() const &
Definition cola_adaptive_index.h:172
bool failed() const noexcept
Definition cola_adaptive_index.h:161
main_pointer main_target() const noexcept
Definition cola_adaptive_index.h:166
auto finish()
Definition cola_adaptive_index.h:169
object_write_paths const & paths() const &&=delete
std::uint64_t spooled_bytes() const noexcept
Definition cola_adaptive_index.h:171
cola_adaptive_index_builder(Factory factory, output_budget budget, std::size_t limit, native_pointer native, main_pointer main={}, native_pointer secondary={})
Definition cola_adaptive_index.h:154
cola_adaptive_index_builder(cola_adaptive_index_builder const &)=delete
bool finished() const noexcept
Definition cola_adaptive_index.h:162
cola_adaptive_index_builder & operator=(cola_adaptive_index_builder const &)=delete
bool spilled() const noexcept
Definition cola_adaptive_index.h:163
Factory const & factory() const noexcept
Definition cola_adaptive_index.h:170
typename builder_type::native_pointer native_pointer
Definition cola_adaptive_index.h:152
std::uint64_t step(std::uint64_t budget)
Definition cola_adaptive_index.h:168
std::uint64_t size() const noexcept
Definition cola_adaptive_index.h:164
native_pointer secondary_target() const noexcept
Definition cola_adaptive_index.h:167
unsigned route
Definition cola_adaptive_index.h:40
adaptive_index_output * owner
Definition cola_adaptive_index.h:39
void append(std::span< std::byte const > bytes)
Definition cola_adaptive_index.h:41
std::uint64_t body_bytes() const noexcept
Definition cola_adaptive_index.h:46
Definition cola_adaptive_index.h:34
adaptive_index_output(adaptive_index_output const &)=delete
Factory const & factory() const noexcept
Definition cola_adaptive_index.h:62
object_write_paths const & paths() const &&=delete
std::array< route_stream, 2 > routes_
Definition cola_adaptive_index.h:113
std::size_t retained_bytes(std::array< bit_string, 2 > const &payload, index_metadata< P > const &metadata) const
Definition cola_adaptive_index.h:129
void poison() noexcept
Definition cola_adaptive_index.h:117
bool failed() const noexcept
Definition cola_adaptive_index.h:59
sink_type main_
Definition cola_adaptive_index.h:114
object_write_paths const & paths() const &
Definition cola_adaptive_index.h:63
result_type finish(typename index_type::native_pointer native, typename index_type::main_pointer main, typename index_type::native_pointer secondary, index_metadata< P > metadata)
Definition cola_adaptive_index.h:77
output_budget budget_
Definition cola_adaptive_index.h:110
adaptive_index_output & operator=(adaptive_index_output const &)=delete
std::size_t limit_
Definition cola_adaptive_index.h:111
std::variant< std::shared_ptr< index_type const >, object_seal_receipt > result_type
Definition cola_adaptive_index.h:37
void require_active() const
Definition cola_adaptive_index.h:128
std::array< borrowed_file_state< P >, 2 > profiles_
Definition cola_adaptive_index.h:115
bool finished() const noexcept
Definition cola_adaptive_index.h:60
bool spilled() const noexcept
Definition cola_adaptive_index.h:61
std::unique_ptr< destination_type > destination_
Definition cola_adaptive_index.h:112
void append_known(unsigned route, bit_view key, std::uint64_t common)
Definition cola_adaptive_index.h:69
destination_type & destination()
Definition cola_adaptive_index.h:121
Factory factory_
Definition cola_adaptive_index.h:109
adaptive_index_output(Factory factory, output_budget budget, std::size_t limit)
Definition cola_adaptive_index.h:52
Definition cola_adaptive_index.h:20
object_stream< P, Ops > stream
Definition cola_adaptive_index.h:22
secondary_spool< SpoolOps > spool
Definition cola_adaptive_index.h:23
index_destination(std::filesystem::path root, object_id id, object_attempt_id attempt, cola_file_dependencies deps, Ops &ops, SpoolOps &spool_ops)
Definition cola_adaptive_index.h:24
cola_file_dependencies dependencies
Definition cola_adaptive_index.h:21
Definition cola_index.h:293
std::array< std::vector< std::byte >, 2 > flags
Definition cola_index.h:295
std::array< rank_groups< P::group_size >, 2 > ranks
Definition cola_index.h:294
std::array< std::vector< std::uint64_t >, 2 > cuts
Definition cola_index.h:296
static index_type adopt(native_pointer native, main_pointer main, native_pointer secondary, std::array< typename index_type::borrowed_array, 2 > borrowed, index_metadata< P > metadata)
Definition cola_index.h:408
Definition cola_file_index.h:39
Definition cola_file_index.h:33
typename index_type::native_pointer native_pointer
Definition cola_index.h:499
Definition cola_index.h:315
std::shared_ptr< native_array const > native_pointer
Definition cola_index.h:322
std::shared_ptr< target_type const > main_pointer
Definition cola_index.h:324
Definition object_writer.h:39
Definition object_path.h:38
Definition object_writer.h:88
Definition object_stream.h:32
void append(std::span< std::byte const > bytes)
Definition object_stream.h:68
Definition object_writer.h:53
Definition output_budget.h:34
Definition output_budget.h:26
static std::shared_ptr< T const > attach(T value, lease allocation)
Definition output_budget.h:83