Everett
Loading...
Searching...
No Matches
cola_file_index.h
Go to the documentation of this file.
1
12#pragma once
13
16
17namespace everett {
18 // Scratch bytes are never objects or durable checkpoints. Unlink the unique
19 // private name immediately after opening it; the descriptor owns its life.
21#if defined(__APPLE__) || defined(__linux__)
22 int create(std::filesystem::path const & path) noexcept {
23 return ::open(path.c_str(), O_RDWR | O_CREAT | O_EXCL | O_CLOEXEC | O_NOFOLLOW, 0600);
24 }
25 int remove(std::filesystem::path const & path) noexcept { return ::unlink(path.c_str()); }
26 std::ptrdiff_t write(int fd, std::span<std::byte const> bytes) noexcept { return ::write(fd, bytes.data(), bytes.size()); }
27 std::ptrdiff_t read_at(int fd, std::span<std::byte> bytes, std::uint64_t at) noexcept {
28 return ::pread(fd, bytes.data(), bytes.size(), static_cast<off_t>(at));
29 }
30 int close(int fd) noexcept { return ::close(fd); }
31#endif
32 };
35 std::optional<blob_identity> main;
36 std::optional<object_id> secondary;
37 };
38 namespace cola_detail {
39 template <class Ops> struct secondary_spool {
40 secondary_spool(std::filesystem::path path, Ops & ops) : path_(std::move(path)), ops_(ops) {}
44 if (fd_ >= 0) (void)ops_.close(fd_);
45 if (named_) (void)ops_.remove(path_);
46 }
47 std::uint64_t body_bytes() const noexcept { return bytes_; }
48 void append(std::span<std::byte const> bytes) {
49 if (bytes.size() > std::uint64_t(std::numeric_limits<std::int64_t>::max()) - bytes_)
50 throw std::length_error("secondary index spool extent");
51 if (bytes.empty()) return;
52 if (fd_ < 0) {
53 fd_ = ops_.create(path_); if (fd_ < 0) fail("create secondary index spool");
54 named_ = true;
55 if (ops_.remove(path_) < 0) fail("unlink secondary index spool");
56 named_ = false;
57 }
58 while (!bytes.empty()) {
59 auto part = bytes.first(std::min(bytes.size(), std::size_t{1} << 20));
60 auto n = ops_.write(fd_, part);
61 if (n < 0 && errno == EINTR) continue;
62 if (n <= 0 || std::uint64_t(n) > part.size()) fail("write secondary index spool", n < 0 ? errno : EIO);
63 bytes_ += std::uint64_t(n); bytes = bytes.subspan(std::size_t(n));
64 }
65 }
66 template <class Stream> void replay(Stream & out) {
67 std::array<std::byte, 64 * 1024> buffer;
68 std::uint64_t at = 0;
69 while (at != bytes_) {
70 auto part = std::span(buffer).first(std::size_t(std::min<std::uint64_t>(bytes_ - at, buffer.size())));
71 auto n = ops_.read_at(fd_, part, at);
72 if (n < 0 && errno == EINTR) continue;
73 if (n <= 0 || std::uint64_t(n) > part.size()) fail("read secondary index spool", n < 0 ? errno : EIO);
74 out.append(std::span<std::byte const>(part).first(std::size_t(n))); at += std::uint64_t(n);
75 }
76 }
77 private:
78 std::filesystem::path path_;
79 Ops & ops_;
80 int fd_ = -1;
81 bool named_ = false;
82 std::uint64_t bytes_ = 0;
83 [[noreturn]] static void fail(char const * operation, int code = errno) {
84 throw std::system_error(code, std::generic_category(), operation);
85 }
86 };
87 template <class P> struct borrowed_file_state {
88 profile_metadata metadata = profile_detail::initial_metadata<P, stream_role::borrowed>();
89 std::vector<std::uint64_t> starts;
91 template <class Sink> void append(Sink & sink, bit_view key, std::uint64_t common) {
92 if (key.size() & (P::bits_per_unit - 1)) throw std::invalid_argument("borrowed file key units");
93 auto retained = common >> P::unit_shift;
94 auto size = key.size() >> P::unit_shift;
95 if (retained > metadata.terminal_key_units || retained > size) throw std::invalid_argument("borrowed file retained prefix");
96 bool block = metadata.record_count % P::codec_block_size == 0;
97 if (block) starts.push_back(sink.position() >> P::unit_shift);
98 auto control = block ? retained : metadata.terminal_key_units - retained;
99 if constexpr (P::unit == profile_unit::byte) count(sink, control);
100 else if (block) sink.template write_count<exponential_golomb<0>>(control);
101 else sink.template write_count<typename P::backspace_encoding>(control);
102 count(sink, size - retained);
103 sink.append(key.subview(retained << P::unit_shift, key.size() - (retained << P::unit_shift)));
105 }
106 void finish_metadata(std::uint64_t bits) {
107 metadata.extent = bits >> P::unit_shift;
108 starts.push_back(metadata.extent); offsets = elias_fano::build<typename P::architecture>(starts);
109 starts.clear(); starts.shrink_to_fit();
110 }
111 template <class Sink> void finish(Sink & sink) {
112 finish_metadata(sink.position()); sink.finish_payload();
113 }
114 private:
115 template <class Sink> static void count(Sink & sink, std::uint64_t value) {
116 if constexpr (P::unit == profile_unit::bit) sink.template write_count<exponential_golomb<0>>(value);
117 else {
118 do { auto byte = unsigned(value & 127); value >>= 7; sink.write_bits(byte | (value ? 128 : 0), 8); } while (value);
119 }
120 }
121 };
122 // Both eager and adaptive outputs finish the same already encoded sections.
123 template <class P, class Sink, class Stream, class Spool>
125 std::array<borrowed_file_state<P>, 2> const & profiles, index_metadata<P> const & metadata,
126 std::uint64_t native_size, Sink & sink, Stream & stream, Spool & spool) {
127 auto emit_offsets = [&](elias_fano const & ef) {
128 sink.align(); sink.words(ef.low); sink.align(); sink.words(ef.high);
129 sink.align(); sink.samples(ef.samples); sink.align(); sink.words(ef.sparse);
130 };
131 std::array<std::byte, cola_section_detail::directory_bytes> directory{};
132 for (unsigned i = 0; i != 4; ++i) directory[i] = std::byte("IX03"[i]);
133 file_detail::put(directory, 4, 2, 3); file_detail::put(directory, 6, 2, cola_section_detail::section_count);
134 file_detail::put(directory, 8, 8, metadata.count); section_detail::put_id(directory, 88, dependencies.native);
135 if (dependencies.main) {
136 directory[82] = std::byte{1}; section_detail::put_id(directory, 104, dependencies.main->native);
137 section_detail::put_id(directory, 120, dependencies.main->index);
138 }
139 if (dependencies.secondary) { directory[83] = std::byte{1}; section_detail::put_id(directory, 136, *dependencies.secondary); }
140 std::array<std::uint64_t, cola_section_detail::section_count> lengths{};
141 std::uint64_t count = 0;
142 for (unsigned route = 0; route != 2; ++route) {
143 auto const & state = profiles[route]; auto const & m = state.metadata; auto const & ef = state.offsets;
144 file_detail::put(directory, 16 + 8 * route, 8, m.record_count);
145 file_detail::put(directory, 32 + 8 * route, 8, m.extent);
146 file_detail::put(directory, 48 + 8 * route, 8, m.terminal_key_units);
147 file_detail::put(directory, 64 + 8 * route, 8, ef.universe); directory[80 + route] = std::byte(ef.low_width);
149 lengths[slot] = profile_detail::byte_count(profile_detail::multiply(m.extent, P::bits_per_unit));
150 lengths[slot + 1] = ef.low.size() * 8; lengths[slot + 2] = ef.high.size() * 8;
151 lengths[slot + 3] = ef.samples.size() * 16; lengths[slot + 4] = ef.sparse.size() * 8;
153 lengths[slot] = metadata.ranks[route].classes.size() * 8;
154 lengths[slot + 1] = metadata.ranks[route].checkpoints.size() * 8;
155 lengths[cola_section_detail::flags_slot(route)] = metadata.flags[route].size();
156 lengths[cola_section_detail::cuts_slot(route)] = metadata.cuts[route].size() * 8;
157 count += m.record_count;
158 }
159 if (count > metadata.count || native_size != metadata.count - count)
160 throw std::invalid_argument("streamed COLA native count");
161 std::uint64_t end = directory.size();
162 for (std::size_t i = 0; i != lengths.size(); ++i) {
163 auto start = profile_detail::add(end, 7) & ~std::uint64_t{7}; end = profile_detail::add(start, lengths[i]);
164 file_detail::put(directory, cola_section_detail::descriptor_offset + (i << 4), 8, start);
165 file_detail::put(directory, cola_section_detail::descriptor_offset + (i << 4) + 8, 8, lengths[i]);
166 }
167 emit_offsets(profiles[0].offsets);
168 sink.align(); spool.replay(stream); emit_offsets(profiles[1].offsets);
169 for (unsigned route = 0; route != 2; ++route) {
170 sink.align(); sink.words(metadata.ranks[route].classes);
171 sink.align(); sink.words(metadata.ranks[route].checkpoints);
172 }
173 for (auto const & flags : metadata.flags) { sink.align(); stream.append(flags); }
174 for (auto const & cuts : metadata.cuts) { sink.align(); sink.words(cuts); }
175 file_header<P> header{file_kind::fractional_index, profile_detail::multiply(end, 1u << (3 - P::unit_shift)), count, 0};
176 return stream.finish(header, directory);
177 }
178 template <class P, class Ops, class SpoolOps> struct file_index_output {
183 file_index_output(std::filesystem::path root, object_id id, object_attempt_id attempt,
184 cola_file_dependencies dependencies, Ops & ops, SpoolOps & spool_ops)
185 : dependencies_(checked(std::move(dependencies))),
186 stream_(std::move(root), std::move(id), std::move(attempt), file_kind::fractional_index,
187 cola_section_detail::directory_bytes, ops),
188 spool_(stream_.paths().private_output.string() + ".secondary", spool_ops), main_(stream_), secondary_(spool_) {}
189 bool failed() const noexcept { return failed_ || stream_.failed(); }
190 bool finished() const noexcept { return stream_.finished(); }
191 object_write_paths const & paths() const & noexcept { return stream_.paths(); }
192 object_write_paths const & paths() const && = delete;
193 std::uint64_t spooled_bytes() const noexcept { return spool_.body_bytes(); }
194 void append_known(unsigned route, bit_view key, std::uint64_t common) {
196 try {
197 if (!route) profiles_[0].append(main_, key, common);
198 else if (route == 1) profiles_[1].append(secondary_, key, common);
199 else throw std::out_of_range("borrowed file route");
200 } catch (...) { failed_ = true; throw; }
201 }
202 template <class Native, class Main> object_seal_receipt finish(std::shared_ptr<Native const> native,
203 std::shared_ptr<Main const> main, std::shared_ptr<Native const> secondary, index_metadata<P> metadata) {
205 try {
206 if (bool(main) != bool(dependencies_.main) || bool(secondary) != bool(dependencies_.secondary) || !native)
207 throw std::invalid_argument("streamed COLA dependency shape");
208 profiles_[0].finish(main_); profiles_[1].finish(secondary_);
209 return seal_file_index<P>(dependencies_, profiles_, metadata, native->size(), main_, stream_, spool_);
210 } catch (...) { failed_ = true; throw; }
211 }
212 private:
218 std::array<borrowed_file_state<P>, 2> profiles_;
219 bool failed_ = false;
220 void require_active() const { if (failed() || finished()) throw std::logic_error("inactive streamed COLA output"); }
222 auto valid = [](object_id const & id) { if (id.hex().size() != 32) throw std::invalid_argument("COLA file dependency identity"); };
223 valid(deps.native); if (deps.main) { valid(deps.main->native); valid(deps.main->index); }
224 if (deps.secondary) valid(*deps.secondary);
225 return deps;
226 }
227 };
228 template <class Output> struct file_index_output_ref {
229 Output * output;
230 bool failed() const noexcept { return output->failed(); }
231 void append_known(unsigned route, bit_view key, std::uint64_t common) { output->append_known(route, key, common); }
232 template <class Native, class Main, class Metadata> auto finish(std::shared_ptr<Native const> native,
233 std::shared_ptr<Main const> main, std::shared_ptr<Native const> secondary, Metadata metadata) {
234 return output->finish(std::move(native), std::move(main), std::move(secondary), std::move(metadata));
235 }
236 };
237 }
238 // Exact sources remain pinned. Only finish seals an adoptable IX03 object;
239 // a partial stream/spool is rebuilt from its source recipe after restart.
240 template <class P, class Native = profile_array<P>, class Main = void,
241 class Ops = posix_object_ops, class SpoolOps = posix_index_spool_ops>
247 cola_file_index_builder(std::filesystem::path root, object_id id, object_attempt_id attempt,
249 : owned_ops_(std::in_place), owned_spool_ops_(std::in_place),
250 output_(std::move(root), std::move(id), std::move(attempt), std::move(dependencies), *owned_ops_, *owned_spool_ops_),
251 builder_({&output_}, std::move(native), std::move(main), std::move(secondary)) {}
252 cola_file_index_builder(std::filesystem::path root, object_id id, object_attempt_id attempt,
254 Ops & ops, SpoolOps & spool_ops)
255 : output_(std::move(root), std::move(id), std::move(attempt), std::move(dependencies), ops, spool_ops),
256 builder_({&output_}, std::move(native), std::move(main), std::move(secondary)) {}
261 bool done() const noexcept { return builder_.done(); }
262 bool failed() const noexcept { return builder_.failed(); }
263 bool finished() const noexcept { return builder_.finished(); }
264 std::uint64_t size() const noexcept { return builder_.size(); }
265 native_pointer native_owner() const noexcept { return builder_.native_owner(); }
266 main_pointer main_target() const noexcept { return builder_.main_target(); }
268 std::uint64_t step(std::uint64_t budget) { return builder_.step(budget); }
270 std::uint64_t spooled_bytes() const noexcept { return output_.spooled_bytes(); }
271 object_write_paths const & paths() const & noexcept { return output_.paths(); }
272 object_write_paths const & paths() const && = delete;
273 private:
274 std::optional<Ops> owned_ops_;
275 std::optional<SpoolOps> owned_spool_ops_;
278 };
279}
Encodes two-target COLA routing in portable IX03 sections.
object_seal_receipt seal_file_index(cola_file_dependencies const &dependencies, std::array< borrowed_file_state< P >, 2 > const &profiles, index_metadata< P > const &metadata, std::uint64_t native_size, Sink &sink, Stream &stream, Spool &spool)
Definition cola_file_index.h:124
unsigned route(unsigned value)
Definition cola_index.h:37
constexpr std::size_t profile_slot(unsigned route) noexcept
Definition cola_sections.h:32
constexpr std::size_t section_count
Definition cola_sections.h:30
constexpr std::size_t flags_slot(unsigned route) noexcept
Definition cola_sections.h:34
constexpr std::size_t cuts_slot(unsigned route) noexcept
Definition cola_sections.h:35
constexpr std::size_t descriptor_offset
Definition cola_sections.h:30
constexpr std::size_t rank_slot(unsigned route) noexcept
Definition cola_sections.h:33
void put(std::span< std::byte > bytes, std::size_t at, unsigned width, std::uint64_t value) noexcept
Definition file.h:70
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
std::uint64_t byte_count(std::uint64_t bits) noexcept
Definition profile.h:50
void put_id(std::span< std::byte > bytes, std::size_t at, object_id const &id)
Definition sections.h:75
Definition active_engine.h:18
file_kind
Definition object_path.h:25
Streams sort-owned KV03 payloads with bounded buffering and shared framing.
Definition profile.h:56
Definition cola_file_index.h:87
profile_metadata metadata
Definition cola_file_index.h:88
elias_fano offsets
Definition cola_file_index.h:90
static void count(Sink &sink, std::uint64_t value)
Definition cola_file_index.h:115
void finish(Sink &sink)
Definition cola_file_index.h:111
std::vector< std::uint64_t > starts
Definition cola_file_index.h:89
void finish_metadata(std::uint64_t bits)
Definition cola_file_index.h:106
void append(Sink &sink, bit_view key, std::uint64_t common)
Definition cola_file_index.h:91
Definition cola_file_index.h:228
bool failed() const noexcept
Definition cola_file_index.h:230
auto finish(std::shared_ptr< Native const > native, std::shared_ptr< Main const > main, std::shared_ptr< Native const > secondary, Metadata metadata)
Definition cola_file_index.h:232
void append_known(unsigned route, bit_view key, std::uint64_t common)
Definition cola_file_index.h:231
Output * output
Definition cola_file_index.h:229
Definition cola_file_index.h:178
object_seal_receipt finish(std::shared_ptr< Native const > native, std::shared_ptr< Main const > main, std::shared_ptr< Native const > secondary, index_metadata< P > metadata)
Definition cola_file_index.h:202
static cola_file_dependencies checked(cola_file_dependencies deps)
Definition cola_file_index.h:221
object_write_paths const & paths() const &&=delete
object_write_paths const & paths() const &noexcept
Definition cola_file_index.h:191
std::array< borrowed_file_state< P >, 2 > profiles_
Definition cola_file_index.h:218
std::uint64_t spooled_bytes() const noexcept
Definition cola_file_index.h:193
cola_file_dependencies dependencies_
Definition cola_file_index.h:213
file_index_output(std::filesystem::path root, object_id id, object_attempt_id attempt, cola_file_dependencies dependencies, Ops &ops, SpoolOps &spool_ops)
Definition cola_file_index.h:183
bool failed() const noexcept
Definition cola_file_index.h:189
bool finished() const noexcept
Definition cola_file_index.h:190
stream_type stream_
Definition cola_file_index.h:214
bool failed_
Definition cola_file_index.h:219
spool_type spool_
Definition cola_file_index.h:215
main_sink main_
Definition cola_file_index.h:216
void require_active() const
Definition cola_file_index.h:220
void append_known(unsigned route, bit_view key, std::uint64_t common)
Definition cola_file_index.h:194
secondary_sink secondary_
Definition cola_file_index.h:217
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
std::uint64_t count
Definition cola_index.h:297
Definition cola_file_index.h:39
std::filesystem::path path_
Definition cola_file_index.h:78
void replay(Stream &out)
Definition cola_file_index.h:66
bool named_
Definition cola_file_index.h:81
std::uint64_t bytes_
Definition cola_file_index.h:82
secondary_spool(std::filesystem::path path, Ops &ops)
Definition cola_file_index.h:40
int fd_
Definition cola_file_index.h:80
Ops & ops_
Definition cola_file_index.h:79
secondary_spool(secondary_spool const &)=delete
std::uint64_t body_bytes() const noexcept
Definition cola_file_index.h:47
static void fail(char const *operation, int code=errno)
Definition cola_file_index.h:83
~secondary_spool()
Definition cola_file_index.h:43
secondary_spool & operator=(secondary_spool const &)=delete
void append(std::span< std::byte const > bytes)
Definition cola_file_index.h:48
Definition cola_file_index.h:33
object_id native
Definition cola_file_index.h:34
std::optional< blob_identity > main
Definition cola_file_index.h:35
std::optional< object_id > secondary
Definition cola_file_index.h:36
Definition cola_file_index.h:242
object_seal_receipt finish()
Definition cola_file_index.h:269
std::uint64_t step(std::uint64_t budget)
Definition cola_file_index.h:268
object_write_paths const & paths() const &&=delete
typename builder_type::main_pointer main_pointer
Definition cola_file_index.h:246
cola_file_index_builder(std::filesystem::path root, object_id id, object_attempt_id attempt, cola_file_dependencies dependencies, native_pointer native, main_pointer main, native_pointer secondary, Ops &ops, SpoolOps &spool_ops)
Definition cola_file_index.h:252
native_pointer native_owner() const noexcept
Definition cola_file_index.h:265
std::uint64_t size() const noexcept
Definition cola_file_index.h:264
cola_file_index_builder & operator=(cola_file_index_builder const &)=delete
object_write_paths const & paths() const &noexcept
Definition cola_file_index.h:271
cola_file_index_builder(std::filesystem::path root, object_id id, object_attempt_id attempt, cola_file_dependencies dependencies, native_pointer native, main_pointer main={}, native_pointer secondary={})
Definition cola_file_index.h:247
std::optional< Ops > owned_ops_
Definition cola_file_index.h:274
native_pointer secondary_target() const noexcept
Definition cola_file_index.h:267
cola_file_index_builder & operator=(cola_file_index_builder &&)=delete
bool done() const noexcept
Definition cola_file_index.h:261
typename builder_type::native_pointer native_pointer
Definition cola_file_index.h:245
cola_file_index_builder(cola_file_index_builder const &)=delete
bool finished() const noexcept
Definition cola_file_index.h:263
bool failed() const noexcept
Definition cola_file_index.h:262
output_type output_
Definition cola_file_index.h:276
main_pointer main_target() const noexcept
Definition cola_file_index.h:266
std::uint64_t spooled_bytes() const noexcept
Definition cola_file_index.h:270
std::optional< SpoolOps > owned_spool_ops_
Definition cola_file_index.h:275
cola_file_index_builder(cola_file_index_builder &&)=delete
builder_type builder_
Definition cola_file_index.h:277
native_pointer secondary_target() const noexcept
Definition cola_index.h:528
std::uint64_t size() const noexcept
Definition cola_index.h:525
main_pointer main_target() const noexcept
Definition cola_index.h:527
bool failed() const noexcept
Definition cola_index.h:523
bool finished() const noexcept
Definition cola_index.h:524
native_pointer native_owner() const noexcept
Definition cola_index.h:526
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 index_type::native_pointer native_pointer
Definition cola_index.h:499
Definition elias_fano.h:397
Definition policy.h:28
Definition file.h:34
Definition object_writer.h:39
Definition object_path.h:38
Definition object_writer.h:88
object_write_paths const & paths() const &noexcept
Definition object_stream.h:62
bool failed() const noexcept
Definition object_stream.h:60
bool finished() const noexcept
Definition object_stream.h:61
Definition object_writer.h:53
Definition cola_file_index.h:20
Definition profile.h:273
std::uint64_t extent
Definition profile.h:291
std::uint64_t terminal_key_units
Definition profile.h:289
std::uint64_t record_count
Definition profile.h:290
Definition sort_profile_file_writer.h:19