16#include <condition_variable>
47 char const *
what() const noexcept
override {
return "session is closed"; }
50 char const *
what() const noexcept
override {
return "queued session contribution was cancelled"; }
77 std::shared_ptr<logical_state const>
logical;
92 bool ready()
const {
return result_.wait_for(std::chrono::seconds(0)) == std::future_status::ready; }
97 std::shared_ptr<identity const>
owner_;
99 std::shared_future<snapshot_type>
result_;
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); });
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>; }
128 return try_submit(Engine::borrow_contribution(std::forward<C>(input)));
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>; }
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>; }
137 template <
class C>
requires std::same_as<std::remove_cvref_t<C>,
contribution_type>
139 return enqueue(std::forward<C>(input),
false);
141 template <
class C>
requires std::same_as<std::remove_cvref_t<C>,
contribution_type>
144 throw std::logic_error(
"blocking submission from the session worker");
145 return *
enqueue(std::forward<C>(input),
true);
147 template <
class C>
requires std::same_as<std::remove_cvref_t<C>,
contribution_type>
154 if (value.owner_ !=
identity_)
return false;
155 auto item = value.request_.lock();
156 if (!item)
return false;
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;
166 std::lock_guard lock(
mutex_);
175 std::lock_guard lock(
mutex_);
179 std::lock_guard lock(
mutex_);
183 std::lock_guard lock(
mutex_);
192 std::lock_guard lock(
mutex_);
199 throw std::logic_error(
"joining the session worker from itself");
207 std::optional<contribution_type>
input;
211 :
input(std::in_place, std::forward<C>(value)),
charge(reservation) {}
222 template <
class C> std::optional<ticket>
enqueue(C && input,
bool wait) {
223 auto charge = Engine::reservation(input);
225 throw std::length_error(
"contribution exceeds a session admission limit");
229 std::unique_lock lock(
mutex_);
233 if (!
fits(charge))
return std::nullopt;
238 std::shared_ptr<request> item;
241 item = std::make_shared<request>(std::forward<C>(input), charge);
244 result.
result_ = item->result.get_future().share();
246 std::lock_guard lock(
mutex_);
255 std::lock_guard lock(
mutex_);
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)});
278 void failed(std::shared_ptr<request>
const &
active, std::exception_ptr error)
noexcept {
279 std::list<std::shared_ptr<request>> abandoned;
281 std::lock_guard lock(
mutex_);
288 auto reject = [&](std::shared_ptr<request>
const & item) {
290 std::lock_guard lock(
mutex_);
291 item->result.set_exception(error);
295 for (
auto const & item : abandoned) reject(item);
299 std::unique_lock lock(
mutex_);
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; }
309 std::shared_ptr<request> item;
312 if constexpr (
requires (Engine
const & value) { { value.admission_ready() } -> std::convertible_to<bool>; }) {
314 if (!
ready && !pending)
throw std::logic_error(
"Engine refuses admission without pending service");
317 std::unique_lock lock(
mutex_);
323 std::optional<world_type> produced;
324 try { produced.emplace(
engine_->contribute(std::move(*item->input))); }
329 if constexpr (std::is_nothrow_move_constructible_v<world_type> &&
330 requires (Engine
const & value) { { value.failed() }
noexcept -> std::same_as<bool>; }) {
332 auto error = std::current_exception();
336 std::lock_guard lock(
mutex_);
337 item->result.set_exception(error);
unreserve(item->charge);
345 auto candidate =
next(std::move(*produced),
true);
348 std::atomic_store_explicit(&
current_, candidate, std::memory_order_release);
350 std::lock_guard lock(
mutex_);
351 item->result.set_value(std::move(candidate));
357 if (candidate) std::atomic_store_explicit(&
current_,
next(std::move(*candidate),
false), std::memory_order_release);
361 failed(item, std::current_exception());
370 std::shared_ptr<identity const>
identity_ = std::make_shared<identity const>();
375 std::list<std::shared_ptr<request>>
queue_;
Definition active_engine.h:18
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
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
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
char const * what() const noexcept override
Definition session.h:50
char const * what() const noexcept override
Definition session.h:47
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
std::uint64_t bytes
Definition session.h:35
std::uint64_t work
Definition session.h:34
bool operator==(session_reservation const &) const =default
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