Everett
Loading...
Searching...
No Matches
index_pipeline_detail.h
Go to the documentation of this file.
1
13#pragma once
14
16
18
19#include <cstddef>
20#include <cstdint>
21#include <memory>
22#include <span>
23#include <stdexcept>
24#include <utility>
25#include <vector>
26
28 // Build indexes from the target outward. Each stage retains one input and
29 // one output sample; sample keys flow directly into the next stage without
30 // materializing an intermediate catalog or rescanning a completed index.
31 // Handoffs carry a policy-unit backspace count and suffix relative to the
32 // preceding sample from that same producer, with a literal first sample.
33 //
34 // A step quantum consumes one augmented occurrence, transfers one sample,
35 // advances the source sampler by at most K occurrences, or communicates EOF.
36 // This is an entry-work budget, not a bound on bytes or wall-clock time.
37 // Each owner separately finalizes its stages and binds exact target identities.
38 template <class P, class Target, class Stage>
40 using policy_type = P;
41 using blob_type = Target;
42 using pair_type = std::shared_ptr<blob_type const>;
43
44 // Stages are stable owners, ordered nearest the source first.
45 pipeline_driver(pair_type target, std::vector<std::unique_ptr<Stage>> stages)
46 : target_(std::move(target)), source_(target_), stages_(std::move(stages)) {
47 emitted_.resize(stages_.size());
48 emitted_suffix_units_.resize(stages_.size());
49 consumed_.resize(stages_.size());
50 }
51
56
57 bool done() const noexcept {
58 return stages_.empty() || stages_.back()->done();
59 }
60 bool finished() const noexcept { return finished_; }
61 bool failed() const noexcept { return failed_; }
62 std::size_t size() const noexcept { return stages_.size(); }
63
64 std::uint64_t step(std::uint64_t quanta) {
65 if (finished_ || failed_) error_detail::raise<std::logic_error>("pipeline is no longer active");
66 std::uint64_t work = 0;
67 try {
68 while (work != quanta && !done()) {
69 if (!advance_one()) error_detail::raise<std::logic_error>("index pipeline made no progress");
70 ++work;
71 }
72 } catch (...) {
73 failed_ = true;
74 throw;
75 }
76 return work;
77 }
78
79 std::span<std::uint64_t const> emitted_samples() const & noexcept { return emitted_; }
80 std::span<std::uint64_t const> emitted_samples() const && = delete;
81 std::span<std::uint64_t const> consumed_entries() const & noexcept { return consumed_; }
82 std::span<std::uint64_t const> consumed_entries() const && = delete;
83 std::span<std::uint64_t const> emitted_suffix_units() const & noexcept { return emitted_suffix_units_; }
84 std::span<std::uint64_t const> emitted_suffix_units() const && = delete;
85 std::uint64_t source_suffix_units() const noexcept { return source_suffix_units_; }
86 sampling_work const & source_work() const & noexcept { return source_.counters(); }
87 sampling_work const & source_work() const && = delete;
88
89 protected:
93 std::vector<std::unique_ptr<Stage>> stages_;
94 std::vector<std::uint64_t> emitted_;
95 std::vector<std::uint64_t> consumed_;
96 std::vector<std::uint64_t> emitted_suffix_units_;
97 std::uint64_t source_suffix_units_ = 0;
98 bool finished_ = false;
99 bool failed_ = false;
100
101 private:
102 bool advance_one() {
103 // Give consumers the first opportunity to release backpressure.
104 for (std::size_t end = stages_.size(); end; --end) {
105 auto i = end - 1;
106 auto & stage = *stages_[i];
107 if (stage.has_output()) {
108 if (i + 1 == stages_.size()) {
109 auto sample = stage.take_coded_output();
111 emitted_suffix_units_[i], (sample.suffix.bit_size >> P::unit_shift));
112 } else {
113 auto & next = *stages_[i + 1];
114 if (!next.needs_input()) continue;
115 auto sample = stage.take_coded_output();
116 next.push(sample);
118 emitted_suffix_units_[i], (sample.suffix.bit_size >> P::unit_shift));
119 }
120 ++emitted_[i];
121 return true;
122 }
123 if (stage.done()) continue;
124 if (stage.needs_input()) {
125 if (!i) {
126 if (source_.done()) {
127 stage.close_input();
128 } else {
129 auto sample = source_.peek();
130 auto coded = source_encoder_.encode(sample.key, sample.target_ordinal);
131 stage.push(coded);
133 source_suffix_units_, (coded.suffix.bit_size >> P::unit_shift));
135 }
136 return true;
137 }
138 if (stages_[i - 1]->done()) {
139 stage.close_input();
140 return true;
141 }
142 continue;
143 }
144 auto consumed = stage.step(1);
145 if (consumed) {
146 consumed_[i] += consumed;
147 return true;
148 }
149 }
150 return false;
151 }
152 };
153}
Outlines exceptional check failures while preserving their types and messages.
Declares Everett's incremental fractional-index builder.
Definition index_pipeline_detail.h:27
std::uint64_t add(std::uint64_t a, std::uint64_t b)
Definition profile.h:39
Definition index_pipeline_detail.h:39
bool finished() const noexcept
Definition index_pipeline_detail.h:60
std::span< std::uint64_t const > emitted_suffix_units() const &noexcept
Definition index_pipeline_detail.h:83
std::vector< std::unique_ptr< Stage > > stages_
Definition index_pipeline_detail.h:93
pipeline_driver & operator=(pipeline_driver &&)=default
bool finished_
Definition index_pipeline_detail.h:98
std::span< std::uint64_t const > emitted_samples() const &&=delete
std::vector< std::uint64_t > emitted_
Definition index_pipeline_detail.h:94
std::uint64_t step(std::uint64_t quanta)
Definition index_pipeline_detail.h:64
P policy_type
Definition index_pipeline_detail.h:40
bool advance_one()
Definition index_pipeline_detail.h:102
std::shared_ptr< blob_type const > pair_type
Definition index_pipeline_detail.h:42
pair_type target_
Definition index_pipeline_detail.h:90
std::span< std::uint64_t const > emitted_suffix_units() const &&=delete
sampling_work const & source_work() const &&=delete
sample_cursor< P, Target > source_
Definition index_pipeline_detail.h:91
profile_sample_encoder< P > source_encoder_
Definition index_pipeline_detail.h:92
pipeline_driver & operator=(pipeline_driver const &)=delete
std::span< std::uint64_t const > consumed_entries() const &noexcept
Definition index_pipeline_detail.h:81
sampling_work const & source_work() const &noexcept
Definition index_pipeline_detail.h:86
std::vector< std::uint64_t > emitted_suffix_units_
Definition index_pipeline_detail.h:96
std::size_t size() const noexcept
Definition index_pipeline_detail.h:62
pipeline_driver(pipeline_driver &&)=default
bool done() const noexcept
Definition index_pipeline_detail.h:57
std::vector< std::uint64_t > consumed_
Definition index_pipeline_detail.h:95
pipeline_driver(pair_type target, std::vector< std::unique_ptr< Stage > > stages)
Definition index_pipeline_detail.h:45
std::uint64_t source_suffix_units() const noexcept
Definition index_pipeline_detail.h:85
bool failed() const noexcept
Definition index_pipeline_detail.h:61
std::span< std::uint64_t const > consumed_entries() const &&=delete
pipeline_driver(pipeline_driver const &)=delete
bool failed_
Definition index_pipeline_detail.h:99
Target blob_type
Definition index_pipeline_detail.h:41
std::uint64_t source_suffix_units_
Definition index_pipeline_detail.h:97
std::span< std::uint64_t const > emitted_samples() const &noexcept
Definition index_pipeline_detail.h:79
Definition sampling.h:82
Definition sampling.h:186
bool done() const noexcept
Definition sampling.h:195
profile_sample_view< P > peek() const &
Definition sampling.h:197
sampling_work const & counters() const noexcept
Definition sampling.h:224
void advance()
Definition sampling.h:205
Definition sampling.h:167