Everett
Loading...
Searching...
No Matches
cola_local_merge.h
Go to the documentation of this file.
1
13#pragma once
14
15#include <everett/cola_index.h>
17
18#include <cstdint>
19#include <memory>
20#include <optional>
21#include <stdexcept>
22#include <utility>
23
24namespace everett {
27
28 // A local destination plan pins exact existing objects. A new main retains
29 // the next level's targets; a new secondary joins an existing main and is a
30 // native-only leaf. This plan assigns no logical slots or admission mass.
31 template <class P> struct cola_destination_plan {
32 using policy_type = P;
36 if (secondary && !main)
37 error_detail::raise<std::invalid_argument>("COLA destination secondary requires a main");
38 return {cola_destination_kind::main, std::move(main), std::move(secondary)};
39 }
41 if (!main) error_detail::raise<std::invalid_argument>("COLA secondary destination requires an existing main");
42 return {cola_destination_kind::secondary, std::move(main), {}};
43 }
44 cola_destination_kind kind() const noexcept { return kind_; }
45 pair_type main_target() const noexcept { return main_; }
46 native_pointer secondary_target() const noexcept { return secondary_; }
47 private:
53 };
54
64
65 // One local COLA job, with explicit native/index/carrier completion. Inputs
66 // must represent adjacent chronological history in older,newer order; sorted
67 // keys cannot establish that semantic premise. No key is elided.
68 //
69 // step counts distinct native keys or augmented index occurrences in only
70 // the current stage. It never performs finish_stage implicitly. Construction
71 // and finish_stage may decode initial heads, allocate, and finalize EF/rank;
72 // their work is not bounded by the record budget. A secondary has no index
73 // stage. The ready carrier is a local empty-native routing artifact, not
74 // necessarily a prepared query root. No slots, visibility, publication,
75 // durable progress or scheduler service guarantee are implemented here.
76 //
77 // Source and plan owners remain pinned for the job's lifetime, including
78 // after failure. A failed execution/finalization poisons this job; a rejected
79 // precondition does not. Moves transfer stable builder objects without
80 // moving user callback state. Assignment releases the previous job's pins.
81 template <class P, class Compose = replace_native_value> struct cola_local_merge_job {
82 using policy_type = P;
91
93 Compose compose = {}, std::optional<std::uint64_t> common = P::value_width)
94 : older_(checked(std::move(older))), newer_(checked(std::move(newer))), plan_(checked(std::move(plan))),
95 merge_(std::make_unique<merge_type>(older_, newer_, std::move(compose), common)) {}
99 cola_local_merge_job & operator=(cola_local_merge_job &&) noexcept = default;
100
102 bool failed() const noexcept { return failed_; }
103 bool finished() const noexcept { return older_ && phase_ == cola_merge_phase::taken; }
104 bool done() const noexcept { return older_ && !failed_ && phase_ == cola_merge_phase::ready; }
105 bool stage_done() const noexcept {
106 if (!older_ || failed_ || finished()) return false;
107 if (phase_ == cola_merge_phase::native_merge) return merge_->done();
109 return index_->done();
111 }
112 std::uint64_t step(std::uint64_t budget = 1) {
114 try {
115 if (phase_ == cola_merge_phase::native_merge) return merge_->step(budget).keys;
117 return index_->step(budget);
118 return 0;
119 } catch (...) { failed_ = true; throw; }
120 }
124 error_detail::raise<std::logic_error>("COLA local stage is not ready for finalization");
125 try {
127 merged_ = std::make_shared<native_type const>(merge_->finish());
128 merge_.reset();
130 index_ = std::make_unique<index_builder_type>(merged_, plan_.main_target(), plan_.secondary_target());
132 } else {
136 }
138 main_ = std::make_shared<index_type const>(index_->finish());
139 index_.reset();
141 } else {
142 carrier_ = std::make_shared<index_type const>(index_->finish());
143 index_.reset();
145 }
146 } catch (...) { failed_ = true; throw; }
147 }
150 if (!done()) error_detail::raise<std::logic_error>("COLA local merge is not ready");
152 return {std::move(merged_), std::move(main_), std::move(secondary_), std::move(carrier_)};
153 }
154 private:
157 std::unique_ptr<merge_type> merge_;
158 std::unique_ptr<index_builder_type> index_;
162 bool failed_ = false;
163
165 if (!source) error_detail::raise<std::invalid_argument>("null COLA local merge source");
166 return source;
167 }
169 if ((!plan.main_target() && plan.secondary_target()) ||
171 error_detail::raise<std::invalid_argument>("invalid COLA destination plan");
172 return plan;
173 }
174 void require_active() const {
175 if (!older_ || failed_ || finished())
176 error_detail::raise<std::logic_error>("COLA local merge is no longer active");
177 }
179 auto empty = std::make_shared<native_type const>(native_type::build({}));
180 index_ = std::make_unique<index_builder_type>(std::move(empty), main_, secondary_);
182 }
183 };
184}
Declares dual-target main/secondary fractional indexes for COLA.
Definition active_engine.h:18
cola_destination_kind
Definition cola_local_merge.h:25
cola_merge_phase
Definition cola_local_merge.h:26
Merges ordered native streams incrementally with policy-specific value composition.
Definition cola_local_merge.h:31
static cola_destination_plan for_secondary(pair_type main)
Definition cola_local_merge.h:40
cola_destination_plan(cola_destination_kind kind, pair_type main, native_pointer secondary)
Definition cola_local_merge.h:51
native_pointer secondary_target() const noexcept
Definition cola_local_merge.h:46
pair_type main_target() const noexcept
Definition cola_local_merge.h:45
native_pointer secondary_
Definition cola_local_merge.h:50
pair_type main_
Definition cola_local_merge.h:49
cola_destination_kind kind_
Definition cola_local_merge.h:48
cola_destination_kind kind() const noexcept
Definition cola_local_merge.h:44
static cola_destination_plan for_main(pair_type main={}, native_pointer secondary={})
Definition cola_local_merge.h:35
typename cola_index< P >::native_pointer native_pointer
Definition cola_local_merge.h:34
P policy_type
Definition cola_local_merge.h:32
typename cola_index< P >::pair_type pair_type
Definition cola_local_merge.h:33
Definition cola_index.h:496
std::shared_ptr< native_array const > native_pointer
Definition cola_index.h:322
std::shared_ptr< cola_index const > pair_type
Definition cola_index.h:323
Definition cola_local_merge.h:81
std::uint64_t step(std::uint64_t budget=1)
Definition cola_local_merge.h:112
bool failed_
Definition cola_local_merge.h:162
P policy_type
Definition cola_local_merge.h:82
void start_carrier()
Definition cola_local_merge.h:178
cola_local_merge_job & operator=(cola_local_merge_job const &)=delete
typename index_type::pair_type pair_type
Definition cola_local_merge.h:86
cola_merge_phase phase_
Definition cola_local_merge.h:161
pair_type main_
Definition cola_local_merge.h:160
void finish_stage()
Definition cola_local_merge.h:121
cola_local_merge_job(cola_local_merge_job const &)=delete
native_merge_builder< P, native_type, Compose > merge_type
Definition cola_local_merge.h:89
native_pointer secondary_
Definition cola_local_merge.h:159
native_pointer newer_
Definition cola_local_merge.h:155
cola_merge_phase phase() const noexcept
Definition cola_local_merge.h:101
native_pointer merged_
Definition cola_local_merge.h:159
static native_pointer checked(native_pointer source)
Definition cola_local_merge.h:164
bool failed() const noexcept
Definition cola_local_merge.h:102
typename index_type::native_pointer native_pointer
Definition cola_local_merge.h:85
std::unique_ptr< merge_type > merge_
Definition cola_local_merge.h:157
void require_active() const
Definition cola_local_merge.h:174
pair_type carrier_
Definition cola_local_merge.h:160
result_type finish()
Definition cola_local_merge.h:148
bool done() const noexcept
Definition cola_local_merge.h:104
cola_local_merge_job(native_pointer older, native_pointer newer, plan_type plan, Compose compose={}, std::optional< std::uint64_t > common=P::value_width)
Definition cola_local_merge.h:92
cola_local_merge_job(cola_local_merge_job &&) noexcept=default
plan_type plan_
Definition cola_local_merge.h:156
static plan_type checked(plan_type plan)
Definition cola_local_merge.h:168
std::unique_ptr< index_builder_type > index_
Definition cola_local_merge.h:158
bool finished() const noexcept
Definition cola_local_merge.h:103
bool stage_done() const noexcept
Definition cola_local_merge.h:105
native_pointer older_
Definition cola_local_merge.h:155
Definition cola_local_merge.h:55
P policy_type
Definition cola_local_merge.h:56
pair_type carrier
Definition cola_local_merge.h:62
native_pointer secondary
Definition cola_local_merge.h:61
typename cola_index< P >::pair_type pair_type
Definition cola_local_merge.h:57
native_pointer merged_native
Definition cola_local_merge.h:59
typename cola_index< P >::native_pointer native_pointer
Definition cola_local_merge.h:58
pair_type main
Definition cola_local_merge.h:60
Definition native_merge.h:174
Definition profile.h:1168
static profile_array build(std::span< profile_record const > records, std::span< std::uint64_t const > prefix_ceilings={}, std::uint64_t restart_factor=0)
Definition profile.h:1173