Everett
Loading...
Searching...
No Matches
session.h
Go to the documentation of this file.
1
13#pragma once
14
15#include <atomic>
16#include <condition_variable>
17#include <chrono>
18#include <concepts>
19#include <cstdint>
20#include <list>
21#include <exception>
22#include <future>
23#include <limits>
24#include <memory>
25#include <mutex>
26#include <optional>
27#include <stdexcept>
28#include <thread>
29#include <type_traits>
30#include <utility>
31
32namespace everett {
34 std::uint64_t work = 0;
35 std::uint64_t bytes = 0;
36 bool operator==(session_reservation const &) const = default;
37 };
38
40 std::uint64_t work;
41 std::uint64_t bytes;
42 std::uint64_t contributions;
43 std::uint64_t maintenance_budget = 128;
44 };
45
46 struct session_closed : std::exception {
47 char const * what() const noexcept override { return "session is closed"; }
48 };
49 struct session_cancelled : std::exception {
50 char const * what() const noexcept override { return "queued session contribution was cancelled"; }
51 };
52
53 // Engine supplies world_type and contribution_type, plus:
54 // static session_reservation reservation(contribution_type const &);
55 // world_type snapshot() const;
56 // world_type contribute(contribution_type);
57 // bool pending() const;
58 // std::optional<world_type> advance(std::uint64_t budget);
59 // reservation is thread-safe and independent of mutable Engine state. Its
60 // work is the Engine's conservative admission allowance. Existing merge
61 // debt is separate; an Engine can expose admission_ready() to require its
62 // service before claiming another queued input.
63 // advance returns only completed, equivalent layouts of the current world.
64 // All returned worlds own their transitive pins independently of Engine.
65 //
66 // One worker owns Engine. Admission limits cover accepted input reservations,
67 // including the running contribution, until its charged path completes. They
68 // do not bound saved snapshots, Engine's working set, total I/O or elapsed time.
69 // shutdown drains accepted contributions but stops optional maintenance. No
70 // user Engine operation or contribution destructor runs under the queue lock.
71 template <class Engine> struct session {
72 using world_type = typename Engine::world_type;
73 using contribution_type = typename Engine::contribution_type;
74
75 struct logical_state {};
76 struct publication {
77 std::shared_ptr<logical_state const> logical;
78 std::uint64_t generation;
79 std::uint64_t revision;
81 };
82 using snapshot_type = std::shared_ptr<publication const>;
83
84 private:
85 struct request;
86 struct identity {};
87
88 public:
89 struct ticket {
90 ticket() = default;
91 bool valid() const noexcept { return result_.valid(); }
92 bool ready() const { return result_.wait_for(std::chrono::seconds(0)) == std::future_status::ready; }
93 void wait() const { result_.wait(); }
94 snapshot_type get() const { return result_.get(); }
95 private:
96 friend struct session;
97 std::shared_ptr<identity const> owner_;
98 std::weak_ptr<request> request_;
99 std::shared_future<snapshot_type> result_;
100 };
101
102 explicit session(Engine engine, session_limits limits)
103 : engine_(std::make_unique<Engine>(std::move(engine))), limits_(limits) {
105 throw std::invalid_argument("session needs a positive contribution limit and maintenance budget");
106 auto first = std::make_shared<publication const>(publication{
107 std::make_shared<logical_state const>(), 0, 0, engine_->snapshot()});
108 std::atomic_store_explicit(&current_, std::move(first), std::memory_order_release);
109 auto pending = engine_->pending();
110 worker_ = std::thread([this, pending] { run(pending); });
111 }
112 session(session const &) = delete;
113 session & operator=(session const &) = delete;
114 session(session &&) = delete;
115 session & operator=(session &&) = delete;
117
118 snapshot_type snapshot() const noexcept { return std::atomic_load_explicit(&current_, std::memory_order_acquire); }
119
120 // A saturated try_submit neither copies nor moves input. Blocking submit
121 // retains no private input copy while waiting. Caller-owned waiting inputs
122 // lie outside the accepted-input reservation bound.
123 // An Engine may adapt a public command to its queue envelope without
124 // taking ownership yet. The ordinary enqueue path still reserves first.
125 template <class C> requires (!std::same_as<std::remove_cvref_t<C>, contribution_type>) &&
126 requires(C && input) { { Engine::borrow_contribution(std::forward<C>(input)) } -> std::same_as<contribution_type>; }
127 std::optional<ticket> try_submit(C && input) {
128 return try_submit(Engine::borrow_contribution(std::forward<C>(input)));
129 }
130 template <class C> requires (!std::same_as<std::remove_cvref_t<C>, contribution_type>) &&
131 requires(C && input) { { Engine::borrow_contribution(std::forward<C>(input)) } -> std::same_as<contribution_type>; }
132 ticket submit(C && input) { return submit(Engine::borrow_contribution(std::forward<C>(input))); }
133 template <class C> requires (!std::same_as<std::remove_cvref_t<C>, contribution_type>) &&
134 requires(C && input) { { Engine::borrow_contribution(std::forward<C>(input)) } -> std::same_as<contribution_type>; }
135 snapshot_type apply(C && input) { return submit(std::forward<C>(input)).get(); }
136
137 template <class C> requires std::same_as<std::remove_cvref_t<C>, contribution_type>
138 std::optional<ticket> try_submit(C && input) {
139 return enqueue(std::forward<C>(input), false);
140 }
141 template <class C> requires std::same_as<std::remove_cvref_t<C>, contribution_type>
142 ticket submit(C && input) {
143 if (worker_session_ == this)
144 throw std::logic_error("blocking submission from the session worker");
145 return *enqueue(std::forward<C>(input), true);
146 }
147 template <class C> requires std::same_as<std::remove_cvref_t<C>, contribution_type>
148 snapshot_type apply(C && input) { return submit(std::forward<C>(input)).get(); }
149
150 // Only an unclaimed queue entry can be cancelled. The ticket remains valid
151 // and resolves with session_cancelled. Waiting, dropping or timing out a ticket
152 // does not cancel it. Other sessions' tickets cannot cancel this queue.
153 bool cancel(ticket const & value) {
154 if (value.owner_ != identity_) return false;
155 auto item = value.request_.lock();
156 if (!item) return false;
157 {
158 std::lock_guard lock(mutex_);
159 auto found = queue_.begin();
160 while (found != queue_.end() && *found != item) ++found;
161 if (found == queue_.end()) return false;
162 queue_.erase(found);
163 }
164 item->input.reset();
165 {
166 std::lock_guard lock(mutex_);
167 item->result.set_exception(std::make_exception_ptr(session_cancelled{}));
168 unreserve(item->charge);
169 }
170 changed_.notify_all();
171 return true;
172 }
173
175 std::lock_guard lock(mutex_);
176 return reserved_;
177 }
178 std::uint64_t pending_count() const {
179 std::lock_guard lock(mutex_);
180 return count_;
181 }
182 std::exception_ptr failure() const {
183 std::lock_guard lock(mutex_);
184 return failure_;
185 }
186
187 // close rejects new admissions, wakes blocked submitters and drains the
188 // accepted queue asynchronously. shutdown additionally joins the worker;
189 // it is idempotent and concurrent shutdown calls serialize their joins.
190 void close() {
191 {
192 std::lock_guard lock(mutex_);
193 closing_ = true;
194 }
195 changed_.notify_all();
196 }
197 void shutdown() {
198 if (worker_session_ == this)
199 throw std::logic_error("joining the session worker from itself");
200 close();
201 std::lock_guard lock(join_mutex_);
202 if (worker_.joinable()) worker_.join();
203 }
204
205 private:
206 struct request {
207 std::optional<contribution_type> input;
209 std::promise<snapshot_type> result;
210 template <class C> request(C && value, session_reservation reservation)
211 : input(std::in_place, std::forward<C>(value)), charge(reservation) {}
212 };
213
214 bool fits(session_reservation charge) const noexcept {
215 return count_ < limits_.contributions && charge.work <= limits_.work - reserved_.work &&
216 charge.bytes <= limits_.bytes - reserved_.bytes;
217 }
218 void available() const {
219 if (failure_) std::rethrow_exception(failure_);
220 if (closing_) throw session_closed{};
221 }
222 template <class C> std::optional<ticket> enqueue(C && input, bool wait) {
223 auto charge = Engine::reservation(input);
224 if (charge.work > limits_.work || charge.bytes > limits_.bytes)
225 throw std::length_error("contribution exceeds a session admission limit");
226 // Reserve before constructing the private input. Copy/move constructors
227 // may run arbitrary code, so reserve under the lock and construct outside.
228 {
229 std::unique_lock lock(mutex_);
230 available();
231 if (wait) changed_.wait(lock, [&] { return closing_ || failure_ || fits(charge); });
232 available();
233 if (!fits(charge)) return std::nullopt;
234 reserved_.work += charge.work;
235 reserved_.bytes += charge.bytes;
236 ++count_;
237 }
238 std::shared_ptr<request> item;
239 ticket result;
240 try {
241 item = std::make_shared<request>(std::forward<C>(input), charge);
242 result.owner_ = identity_;
243 result.request_ = item;
244 result.result_ = item->result.get_future().share();
245 {
246 std::lock_guard lock(mutex_);
247 // close drains this reservation too; a failure cannot accept a new
248 // request after the failed worker has already emptied its queue.
249 if (failure_) std::rethrow_exception(failure_);
250 queue_.push_back(item);
251 }
252 } catch (...) {
253 item.reset();
254 {
255 std::lock_guard lock(mutex_);
256 unreserve(charge);
257 }
258 changed_.notify_all();
259 throw;
260 }
261 changed_.notify_all();
262 return result;
263 }
264 void unreserve(session_reservation charge) noexcept {
265 reserved_.work -= charge.work;
266 reserved_.bytes -= charge.bytes;
267 --count_;
268 }
269 snapshot_type next(world_type world, bool contribution) const {
270 auto old = snapshot();
271 if (old->revision == std::numeric_limits<std::uint64_t>::max() ||
272 (contribution && old->generation == std::numeric_limits<std::uint64_t>::max()))
273 throw std::length_error("session publication sequence exhausted");
274 return std::make_shared<publication const>(publication{
275 contribution ? std::make_shared<logical_state const>() : old->logical,
276 old->generation + std::uint64_t(contribution), old->revision + 1, std::move(world)});
277 }
278 void failed(std::shared_ptr<request> const & active, std::exception_ptr error) noexcept {
279 std::list<std::shared_ptr<request>> abandoned;
280 {
281 std::lock_guard lock(mutex_);
282 failure_ = error;
283 closing_ = true;
284 abandoned.swap(queue_);
285 }
286 // Destroy failed continuations before refunding their input reservation.
287 engine_.reset();
288 auto reject = [&](std::shared_ptr<request> const & item) {
289 item->input.reset();
290 std::lock_guard lock(mutex_);
291 item->result.set_exception(error);
292 unreserve(item->charge);
293 };
294 if (active) reject(active);
295 for (auto const & item : abandoned) reject(item);
296 changed_.notify_all();
297 // Input constructors and a simultaneous cancellation run outside the
298 // lock. Their reservations must also settle before shutdown can join.
299 std::unique_lock lock(mutex_);
300 changed_.wait(lock, [&] { return count_ == 0; });
301 }
302 void run(bool pending) noexcept {
303 struct worker_scope {
304 session const * previous = worker_session_;
305 explicit worker_scope(session const * self) { worker_session_ = self; }
306 ~worker_scope() { worker_session_ = previous; }
307 } scope(this);
308 while (true) {
309 std::shared_ptr<request> item;
310 try {
311 bool ready = true;
312 if constexpr (requires (Engine const & value) { { value.admission_ready() } -> std::convertible_to<bool>; }) {
313 ready = engine_->admission_ready();
314 if (!ready && !pending) throw std::logic_error("Engine refuses admission without pending service");
315 }
316 {
317 std::unique_lock lock(mutex_);
318 changed_.wait(lock, [&] { return !queue_.empty() || (closing_ ? count_ == 0 : pending); });
319 if (!queue_.empty() && ready) { item = std::move(queue_.front()); queue_.pop_front(); }
320 else if (closing_ && count_ == 0) break;
321 }
322 if (item) {
323 std::optional<world_type> produced;
324 try { produced.emplace(engine_->contribute(std::move(*item->input))); }
325 catch (...) {
326 // This optional contract certifies that a rejected contribution
327 // did not change logical state. Do not apply it to failures in
328 // next(), durability publication, or other work after contribute.
329 if constexpr (std::is_nothrow_move_constructible_v<world_type> &&
330 requires (Engine const & value) { { value.failed() } noexcept -> std::same_as<bool>; }) {
331 if (!engine_->failed()) {
332 auto error = std::current_exception();
333 pending = engine_->pending();
334 item->input.reset();
335 {
336 std::lock_guard lock(mutex_);
337 item->result.set_exception(error); unreserve(item->charge);
338 }
339 changed_.notify_all();
340 continue;
341 }
342 }
343 throw;
344 }
345 auto candidate = next(std::move(*produced), true);
346 item->input.reset();
347 pending = engine_->pending();
348 std::atomic_store_explicit(&current_, candidate, std::memory_order_release);
349 {
350 std::lock_guard lock(mutex_);
351 item->result.set_value(std::move(candidate));
352 unreserve(item->charge);
353 }
354 changed_.notify_all();
355 } else {
356 auto candidate = engine_->advance(limits_.maintenance_budget);
357 if (candidate) std::atomic_store_explicit(&current_, next(std::move(*candidate), false), std::memory_order_release);
358 pending = engine_->pending();
359 }
360 } catch (...) {
361 failed(item, std::current_exception());
362 return;
363 }
364 }
365 engine_.reset();
366 }
367
368 std::unique_ptr<Engine> engine_;
370 std::shared_ptr<identity const> identity_ = std::make_shared<identity const>();
372 mutable std::mutex mutex_;
373 std::mutex join_mutex_;
374 std::condition_variable changed_;
375 std::list<std::shared_ptr<request>> queue_;
377 std::uint64_t count_ = 0;
378 bool closing_ = false;
379 std::exception_ptr failure_;
380 std::thread worker_;
381 inline static thread_local session const * worker_session_ = nullptr;
382 };
383}
Definition active_engine.h:18
Definition session.h:86
Definition session.h:75
Definition session.h:76
std::uint64_t revision
Definition session.h:79
std::shared_ptr< logical_state const > logical
Definition session.h:77
std::uint64_t generation
Definition session.h:78
world_type world
Definition session.h:80
Definition session.h:206
std::optional< contribution_type > input
Definition session.h:207
request(C &&value, session_reservation reservation)
Definition session.h:210
session_reservation charge
Definition session.h:208
std::promise< snapshot_type > result
Definition session.h:209
Definition session.h:89
std::weak_ptr< request > request_
Definition session.h:98
bool ready() const
Definition session.h:92
bool valid() const noexcept
Definition session.h:91
snapshot_type get() const
Definition session.h:94
std::shared_future< snapshot_type > result_
Definition session.h:99
void wait() const
Definition session.h:93
std::shared_ptr< identity const > owner_
Definition session.h:97
Definition session.h:49
char const * what() const noexcept override
Definition session.h:50
Definition session.h:46
char const * what() const noexcept override
Definition session.h:47
Definition session.h:39
std::uint64_t work
Definition session.h:40
std::uint64_t bytes
Definition session.h:41
std::uint64_t contributions
Definition session.h:42
std::uint64_t maintenance_budget
Definition session.h:43
Definition session.h:33
std::uint64_t bytes
Definition session.h:35
std::uint64_t work
Definition session.h:34
bool operator==(session_reservation const &) const =default
Definition session.h:71
std::optional< ticket > enqueue(C &&input, bool wait)
Definition session.h:222
bool fits(session_reservation charge) const noexcept
Definition session.h:214
session & operator=(session const &)=delete
session_reservation outstanding() const
Definition session.h:174
std::exception_ptr failure() const
Definition session.h:182
std::shared_ptr< identity const > identity_
Definition session.h:370
void shutdown()
Definition session.h:197
typename Engine::world_type world_type
Definition session.h:72
static thread_local session const * worker_session_
Definition session.h:381
snapshot_type apply(C &&input)
Definition session.h:135
std::mutex join_mutex_
Definition session.h:373
ticket submit(C &&input)
Definition session.h:132
snapshot_type snapshot() const noexcept
Definition session.h:118
session(session &&)=delete
std::unique_ptr< Engine > engine_
Definition session.h:368
~session()
Definition session.h:116
snapshot_type apply(C &&input)
Definition session.h:148
std::optional< ticket > try_submit(C &&input)
Definition session.h:127
session_limits limits_
Definition session.h:369
void failed(std::shared_ptr< request > const &active, std::exception_ptr error) noexcept
Definition session.h:278
snapshot_type next(world_type world, bool contribution) const
Definition session.h:269
std::exception_ptr failure_
Definition session.h:379
std::uint64_t pending_count() const
Definition session.h:178
session(Engine engine, session_limits limits)
Definition session.h:102
std::list< std::shared_ptr< request > > queue_
Definition session.h:375
std::condition_variable changed_
Definition session.h:374
std::mutex mutex_
Definition session.h:372
bool cancel(ticket const &value)
Definition session.h:153
void unreserve(session_reservation charge) noexcept
Definition session.h:264
snapshot_type current_
Definition session.h:371
bool closing_
Definition session.h:378
void close()
Definition session.h:190
session(session const &)=delete
session & operator=(session &&)=delete
session_reservation reserved_
Definition session.h:376
std::uint64_t count_
Definition session.h:377
std::shared_ptr< publication const > snapshot_type
Definition session.h:82
typename Engine::contribution_type contribution_type
Definition session.h:73
void run(bool pending) noexcept
Definition session.h:302
ticket submit(C &&input)
Definition session.h:142
std::thread worker_
Definition session.h:380
void available() const
Definition session.h:218
std::optional< ticket > try_submit(C &&input)
Definition session.h:138
Definition multiverse.h:41