Source
stdlib/memory/owner.hpp
1
// Copyright (c) 2026 BigBrain LLC. MIT-licensed (see LICENSE).2
// Original work; see ACKNOWLEDGMENTS.md for the open-source ideas we build upon.3
#pragma once5
/**6
* @file owner.hpp7
* @brief `memory::Owner<T>` — the sole owner + scheduling coordinator — and the `own()` factory.8
*9
* `T` is the only class template argument; scheduling `policy` is a constructor value, a write's10
* priority is a compile-time argument on `rwrite`. Non-copyable and pinned (so `&value_` is stable11
* forever — a renewed reader never dangles even when the value's internal buffer reallocates).12
*13
* The coordinator is a hand-rolled priority reader/writer lock (one `std::mutex` + `condition_variable`14
* over explicit state — `std::shared_mutex` can't honor write priorities). Guarantees:15
* - readers share; a writer is exclusive (no torn reads, no lost updates) — synchronized via the16
* coordinator mutex at every grant/release, so object access is race-free with no lock held during use.17
* - **drain-before-write**: a write request flips the current read generation's gate, so readers that18
* loop on `valid()` yield; the write proceeds only once `readers_ == 0`.19
* - **priority**: waiting writes are ordered by `(priority desc, arrival asc)` in a priority_queue.20
* - **immediate-write** (`rwrite<memory::immediate>()`, i.e. priority < 0): bypasses the queue; if a21
* writer is active AND cooperating (loops on `valid()`), it preempts that writer, does its write,22
* then the writer resumes; against a non-looping writer it simply waits for it to finish, then goes23
* ahead of the queue.24
*/26
#include <atomic>27
#include <condition_variable>28
#include <cstdint>29
#include <functional>30
#include <future>31
#include <memory>32
#include <mutex>33
#include <queue>34
#include <utility>35
#include <vector>37
#include "lease.hpp"38
#include "mode.hpp"39
#include "ownable.hpp"40
#include "policy.hpp"41
#include "request.hpp"43
namespace cheatah::memory {45
/**46
* @brief The sole owner + scheduling coordinator for one `T`. Non-copyable and pinned, so the object47
* never moves; every access goes through a request → acquire → lease. @tparam T the owned type.48
*/49
template <Ownable T>50
class Owner {51
public:52
/// Take sole ownership by MOVING @p value in — the object is consumed (its resources move into the53
/// Owner), never copied. An lvalue won't bind here (bind an rvalue: `own(std::move(x))`).54
/// @param value the object to take ownership of (moved in).55
/// @param pol the scheduling policy (interleave / writes_first).56
/// @complexity O(1) plus moving @p value.57
/// @alloc one read-generation gate (`std::make_shared`); the value's own storage just moves.58
/// @concurrency construct before sharing; the pinned Owner must outlive every lease and every59
/// thread that uses it.60
/// @test Memory.OwnerConsumesAndMovesTheObjectInNeverCopies61
explicit Owner(T&& value, policy pol = policy::interleave)62
: value_(std::move(value)), policy_(pol), read_gate_(make_gate()) {} // declaration order: policy_ precedes read_gate_, so make_gate() sees it63
Owner(const Owner&) = delete; // copying an owner is forbidden — there is one owner.64
Owner& operator=(const Owner&) = delete; // sole ownership; pinned (mutex is non-movable).65
Owner(Owner&&) = delete; // pinned: &value_ must stay stable for live leases.66
Owner& operator=(Owner&&) = delete;67
~Owner() = default; // the value and the coordinator state unwind in declaration order.69
/// Request a shared READ lease. Blocks — in THIS call: the grant is synchronous, so the returned70
/// request is already fulfilled — only while a write is pending/active or queued; otherwise many71
/// read leases coexist. @return a request for a read lease.72
/// @complexity O(1) plus the blocking wait. @alloc the request's promise/future.73
/// @concurrency callable from any thread; readers share. Writer-preference: this waits while any74
/// write is active, suspended, or queued, so writers cannot starve.75
/// @warning requesting while the SAME thread still holds a lease on this owner can deadlock (a76
/// queued write makes the grant wait on that very lease) — renewal is release, then re-request.77
/// @test Memory.ReadLeasesCoexist78
/// @test MemoryConcurrency.ManyReadersCoexistThenAWriteDrains79
/// @systest MemoryCheatah.ReadLeaseValidState80
Request<Lease<T, read>> rread() {81
std::promise<Lease<T, read>> p;82
auto fut = p.get_future();83
{84
std::unique_lock<std::mutex> lk(mtx_);85
cv_.wait(lk, [&] { return can_read(); });86
++readers_;87
auto gate = read_gate_;88
p.set_value(Lease<T, read>(&value_, std::move(gate), [this] { release_read(); }));89
}90
return Request<Lease<T, read>>(std::move(fut));91
}93
/// Request an exclusive WRITE lease at compile-time `priority` (a plain int or the caller's enum;94
/// higher = served first, ties FIFO). `priority < 0` (spell it `memory::immediate`) is an95
/// immediate-write. Blocks in THIS call — the grant is synchronous, so the returned request is96
/// already fulfilled — until the readers drain and this write wins the queue.97
/// @tparam priority the compile-time write priority. @return a request for a write98
/// lease. @complexity O(log k) to enqueue among k waiters (O(1) immediate), plus the blocking wait.99
/// @alloc the request's promise/future, two fresh gates (`std::make_shared`: this write's own +100
/// the next read generation's), plus one queue-ticket slot (amortized) for a non-immediate write.101
/// @concurrency callable from any thread. Drain-before-write: flips the current read generation's102
/// gate and waits until every reader has released and no other write is active; an immediate-write103
/// skips the queue and additionally preempts a cooperating active writer (which resumes after).104
/// @warning requesting while the SAME thread still holds a lease on this owner deadlocks (the105
/// drain waits on that very lease) — release first, then re-request.106
/// @test Memory.WriteWaitsForReadersToDrain107
/// @test Memory.HigherPriorityWriteServedFirst108
/// @test Memory.NegativePriorityImmediateWritePreemptsTheActiveWriterWhichThenResumes109
/// @test MemoryConcurrency.ManyWritersDeterministicSum110
/// @systest MemoryCheatah.ConcurrentSumOverSharedOwner111
template <auto priority = 0>112
Request<Lease<T, write>> rwrite() {113
constexpr auto P = static_cast<long long>(priority);114
if constexpr (P < 0) return grant_immediate();115
else return grant_write(P);116
}118
private:119
// ── the object + coordinator state (all guarded by mtx_) ──120
T value_;121
policy policy_;122
std::mutex mtx_;123
std::condition_variable cv_;124
long long readers_ = 0; ///< active read leases.125
bool writer_ = false; ///< an active (non-suspended) write lease.126
bool writer_suspended_ = false; ///< active writer paused for an immediate-write.127
bool immediate_ = false;///< an immediate-write holds exclusive access.128
std::shared_ptr<detail::Gate> read_gate_; ///< current read generation's yield gate.129
std::shared_ptr<detail::Gate> writer_gate_; ///< the active writer's gate (for preempt/resume).131
struct Ticket { long long prio; std::uint64_t seq; };132
struct ServedLater { // priority_queue is a max-heap: top = highest priority, then earliest arrival.133
/**134
* The heap ordering: is @p a served later than @p b? Higher priority wins; within a135
* priority, the earlier arrival (lower seq) wins — FIFO among equals.136
* @param a one waiting write's ticket.137
* @param b the other waiting write's ticket.138
* @return true iff @p a is served after @p b.139
* @complexity O(1).140
* @alloc none.141
* @test Memory.HigherPriorityWriteServedFirst142
*/143
bool operator()(const Ticket& a, const Ticket& b) const {144
return a.prio != b.prio ? a.prio < b.prio : a.seq > b.seq;145
}146
};147
std::priority_queue<Ticket, std::vector<Ticket>, ServedLater> wq_; ///< waiting non-immediate writes.148
std::uint64_t seq_ = 0;150
// Readers proceed only when no writer/immediate is active or paused and no writer is queued151
// (writer-preference — this is what forces the drain and prevents writer starvation).152
bool can_read() const { return !writer_ && !immediate_ && !writer_suspended_ && wq_.empty(); }154
// A gate whose holder, on observing !valid, wakes our cv_. The wake must synchronize on155
// mtx_ before notifying: the holder acks from outside the lock, so a bare notify_all could156
// land while grant_immediate() still holds mtx_ evaluating its wait predicate (acked read157
// as false, waiter not yet blocked) — and since the ack is one-shot, that lost wakeup left158
// the preempting write asleep forever against a writer spinning on valid(). Taking and159
// releasing mtx_ first pins the wake after the waiter is actually waiting.160
std::shared_ptr<detail::Gate> make_gate() {161
auto g = std::make_shared<detail::Gate>();162
g->wake = [this] {163
{ std::lock_guard<std::mutex> lk(mtx_); } // serialize with a waiter mid-predicate164
cv_.notify_all();165
};166
return g;167
}169
Request<Lease<T, write>> grant_write(long long prio) {170
std::promise<Lease<T, write>> p;171
auto fut = p.get_future();172
{173
std::unique_lock<std::mutex> lk(mtx_);174
const std::uint64_t my = ++seq_;175
wq_.push({prio, my});176
read_gate_->valid.store(false); // ask current readers to yield (drain)177
cv_.notify_all();178
cv_.wait(lk, [&] {179
return !writer_ && !immediate_ && !writer_suspended_ && readers_ == 0 &&180
!wq_.empty() && wq_.top().seq == my; // no active access + I'm the winner181
});182
wq_.pop();183
writer_ = true;184
writer_gate_ = make_gate(); // fresh, valid — this writer is preemptible185
read_gate_ = make_gate(); // fresh read generation for future readers186
p.set_value(Lease<T, write>(&value_, writer_gate_, [this] { release_write(); }));187
}188
return Request<Lease<T, write>>(std::move(fut));189
}191
Request<Lease<T, write>> grant_immediate() {192
std::promise<Lease<T, write>> p;193
auto fut = p.get_future();194
{195
std::unique_lock<std::mutex> lk(mtx_);196
if (writer_) { // preempt a cooperating active writer197
writer_gate_->valid.store(false);198
cv_.notify_all();199
// Short-circuit !writer_ FIRST: a non-looping writer may release (writer_gate_ ->200
// nullptr) during the wait, so never deref the gate once the writer is gone.201
cv_.wait(lk, [&] { return !writer_ || writer_gate_->acked.load(); });202
if (writer_) { writer_suspended_ = true; writer_ = false; } // it paused → suspend it203
// else it finished on its own; no resume owed.204
}205
read_gate_->valid.store(false); // drain any readers206
cv_.notify_all();207
cv_.wait(lk, [&] { return readers_ == 0 && !immediate_; });208
immediate_ = true;209
read_gate_ = make_gate();210
auto ig = make_gate();211
p.set_value(Lease<T, write>(&value_, std::move(ig), [this] { release_immediate(); }));212
}213
return Request<Lease<T, write>>(std::move(fut));214
}216
void release_read() {217
std::lock_guard<std::mutex> lk(mtx_);218
--readers_;219
cv_.notify_all();220
}221
void release_write() {222
std::lock_guard<std::mutex> lk(mtx_);223
writer_ = false;224
writer_gate_ = nullptr;225
cv_.notify_all();226
}227
void release_immediate() {228
std::lock_guard<std::mutex> lk(mtx_);229
immediate_ = false;230
if (writer_suspended_) { // resume the writer we preempted (before the queue)231
writer_suspended_ = false;232
writer_ = true;233
writer_gate_->acked.store(false);234
writer_gate_->valid.store(true); // its valid() flips back to true → it continues235
}236
cv_.notify_all();237
}238
};240
/// Take sole ownership of @p value and hand back its `Owner`. @tparam T the owned type.241
/// @param value the object to own (moved in). @param pol the scheduling policy.242
/// @return an `Owner<T>` that has consumed @p value.243
/// @complexity O(1) plus moving @p value. @alloc one read-generation gate (`std::make_shared`, in244
/// the `Owner` constructor); the value's own storage just moves.245
/// @test Memory.ObjectDiesWithOwner246
template <Ownable T>247
Owner<T> own(T value, policy pol = policy::interleave) { return Owner<T>(std::move(value), pol); }249
} // namespace cheatah::memory