cheatah
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 once
5/**
6 * @file owner.hpp
7 * @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's
10 * priority is a compile-time argument on `rwrite`. Non-copyable and pinned (so `&value_` is stable
11 * 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 the
16 * 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 that
18 * 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 a
21 * 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 goes
23 * 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"
43namespace cheatah::memory {
45/**
46 * @brief The sole owner + scheduling coordinator for one `T`. Non-copyable and pinned, so the object
47 * never moves; every access goes through a request → acquire → lease. @tparam T the owned type.
48 */
49template <Ownable T>
50class Owner {
51public:
52 /// Take sole ownership by MOVING @p value in — the object is consumed (its resources move into the
53 /// 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 every
59 /// thread that uses it.
60 /// @test Memory.OwnerConsumesAndMovesTheObjectInNeverCopies
61 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 it
63 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 returned
70 /// request is already fulfilled — only while a write is pending/active or queued; otherwise many
71 /// 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 any
74 /// 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 (a
76 /// queued write makes the grant wait on that very lease) — renewal is release, then re-request.
77 /// @test Memory.ReadLeasesCoexist
78 /// @test MemoryConcurrency.ManyReadersCoexistThenAWriteDrains
79 /// @systest MemoryCheatah.ReadLeaseValidState
80 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 an
95 /// immediate-write. Blocks in THIS call — the grant is synchronous, so the returned request is
96 /// 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 write
98 /// 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's
102 /// gate and waits until every reader has released and no other write is active; an immediate-write
103 /// 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 (the
105 /// drain waits on that very lease) — release first, then re-request.
106 /// @test Memory.WriteWaitsForReadersToDrain
107 /// @test Memory.HigherPriorityWriteServedFirst
108 /// @test Memory.NegativePriorityImmediateWritePreemptsTheActiveWriterWhichThenResumes
109 /// @test MemoryConcurrency.ManyWritersDeterministicSum
110 /// @systest MemoryCheatah.ConcurrentSumOverSharedOwner
111 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 }
118private:
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 a
135 * 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.HigherPriorityWriteServedFirst
142 */
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 queued
151 // (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 on
155 // mtx_ before notifying: the holder acks from outside the lock, so a bare notify_all could
156 // land while grant_immediate() still holds mtx_ evaluating its wait predicate (acked read
157 // as false, waiter not yet blocked) — and since the ack is one-shot, that lost wakeup left
158 // the preempting write asleep forever against a writer spinning on valid(). Taking and
159 // 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-predicate
164 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 winner
181 });
182 wq_.pop();
183 writer_ = true;
184 writer_gate_ = make_gate(); // fresh, valid — this writer is preemptible
185 read_gate_ = make_gate(); // fresh read generation for future readers
186 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 writer
197 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 it
203 // else it finished on its own; no resume owed.
204 }
205 read_gate_->valid.store(false); // drain any readers
206 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 continues
235 }
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`, in
244/// the `Owner` constructor); the value's own storage just moves.
245/// @test Memory.ObjectDiesWithOwner
246template <Ownable T>
247Owner<T> own(T value, policy pol = policy::interleave) { return Owner<T>(std::move(value), pol); }
249} // namespace cheatah::memory