The Pedigree Project 0.1
PipeBuffer-transfer.cc
1/* Copyright (c) 2026, Pedigree Developers. */
2#include "pedigree/kernel/utilities/assert.h"
3
4#include "PipeBuffer.h"
5
6PipeBuffer::ReadReservation::ReadReservation() = default;
7
8PipeBuffer::ReadReservation::~ReadReservation() {
9 cancel();
10}
11
12size_t PipeBuffer::ReadReservation::size() const {
13 return m_Size;
14}
15
16void PipeBuffer::ReadReservation::copyTo(uint8_t* destination, size_t count) const {
17 assert(count <= m_Size);
18 if (!count) {
19 return;
20 }
21 assert(m_Buffer);
22 PipeBuffer::lock(m_Buffer->m_Lock);
23 m_Buffer->copyOutLocked(destination, count);
24 m_Buffer->m_Lock.release();
25}
26
27void PipeBuffer::ReadReservation::consume(size_t accepted) {
28 assert(accepted <= m_Size);
29 PipeBuffer* buffer = m_Buffer;
30 if (!buffer) {
31 return;
32 }
33 PipeBuffer::lock(buffer->m_Lock);
34 assert(buffer->m_ReadReserved);
35 buffer->consumeLocked(accepted);
36 buffer->m_ReadReserved = false;
37 buffer->finishResetLocked();
38 buffer->changedLocked();
39 m_Buffer = nullptr;
40 m_Size = 0;
41 buffer->m_Lock.release();
42 buffer->publish();
43 buffer->endOperation();
44}
45
46void PipeBuffer::ReadReservation::cancel() {
47 consume(0);
48}
49
50PipeBuffer::WriteReservation::WriteReservation() = default;
51
52PipeBuffer::WriteReservation::~WriteReservation() {
53 cancel();
54}
55
56size_t PipeBuffer::WriteReservation::size() const {
57 return m_Size;
58}
59
60PipeBuffer::Result PipeBuffer::WriteReservation::commit(const uint8_t* source, size_t accepted) {
61 PipeBuffer* buffer = m_Buffer;
62 if (!buffer) {
63 return {Status::Invalid, 0};
64 }
65 if (accepted > m_Size) {
66 cancel();
67 return {Status::Invalid, 0};
68 }
69 PipeBuffer::lock(buffer->m_Lock);
70 assert(buffer->m_WriteReserved);
71 // A reopen/reset invalidates an old append admission before any source is consumed.
72 const bool closed = buffer->m_Closing || !buffer->m_ReadEnabled || !buffer->m_WriteEnabled ||
73 buffer->m_ResetPending;
74 if (!closed) {
75 buffer->appendLocked(source, accepted);
76 }
77 buffer->m_WriteReserved = false;
78 buffer->finishResetLocked();
79 buffer->changedLocked();
80 m_Buffer = nullptr;
81 m_Size = 0;
82 buffer->m_Lock.release();
83 buffer->publish();
84 buffer->endOperation();
85 return {closed ? Status::Closed : Status::Ready, closed ? 0 : accepted};
86}
87
88void PipeBuffer::WriteReservation::cancel() {
89 PipeBuffer* buffer = m_Buffer;
90 if (!buffer) {
91 return;
92 }
93 PipeBuffer::lock(buffer->m_Lock);
94 assert(buffer->m_WriteReserved);
95 buffer->m_WriteReserved = false;
96 buffer->finishResetLocked();
97 buffer->changedLocked();
98 m_Buffer = nullptr;
99 m_Size = 0;
100 buffer->m_Lock.release();
101 buffer->publish();
102 buffer->endOperation();
103}
104
105PipeBuffer::Result PipeBuffer::reserveRead(size_t maximum, bool block,
106 ReadReservation& reservation) {
107 if (reservation.m_Buffer) {
108 return {Status::Invalid, 0};
109 }
110 if (!maximum) {
111 return {Status::Ready, 0};
112 }
113 ActiveOperation operation(*this);
114 if (!operation) {
115 return {Status::Eof, 0};
116 }
117 lock(m_Lock);
118 Result result = waitLocked(false, 1, block);
119 if (result.status == Status::Ready) {
120 result.count = maximum < result.count ? maximum : result.count;
121 reservation.m_Buffer = this;
122 reservation.m_Size = result.count;
123 m_ReadReserved = true;
124 changedLocked();
125 operation.detach();
126 }
127 m_Lock.release();
128 if (result.status == Status::Ready) {
129 publish();
130 }
131 return result;
132}
133
134PipeBuffer::Result PipeBuffer::reserveWrite(size_t maximum, bool block,
135 WriteReservation& reservation) {
136 if (reservation.m_Buffer) {
137 return {Status::Invalid, 0};
138 }
139 if (!maximum) {
140 return {Status::Ready, 0};
141 }
142 ActiveOperation operation(*this);
143 if (!operation) {
144 return {Status::Closed, 0};
145 }
146 lock(m_Lock);
147 Result result = waitLocked(true, 1, block);
148 if (result.status == Status::Ready) {
149 result.count = maximum < result.count ? maximum : result.count;
150 reservation.m_Buffer = this;
151 reservation.m_Size = result.count;
152 m_WriteReserved = true;
153 changedLocked();
154 operation.detach();
155 }
156 m_Lock.release();
157 if (result.status == Status::Ready) {
158 publish();
159 }
160 return result;
161}
162
164 Mutex mutex;
165 ConditionVariable condition;
166 bool pending = false;
167};
168
171 PairLink* next = nullptr;
172};
173
174void PipeBuffer::linkPairLocked(PairLink& link) {
175 link.next = m_PairWaiters;
176 m_PairWaiters = &link;
177}
178
179void PipeBuffer::unlinkPairLocked(PairLink& link) {
180 PairLink** position = &m_PairWaiters;
181 while (*position && *position != &link) {
182 position = &(*position)->next;
183 }
184 assert(*position == &link);
185 *position = link.next;
186 link.next = nullptr;
187}
188
189void PipeBuffer::wakePairsLocked() {
190 for (PairLink* link = m_PairWaiters; link; link = link->next) {
191 lock(link->waiter->mutex);
192 link->waiter->pending = true;
193 link->waiter->condition.signal();
194 link->waiter->mutex.release();
195 }
196}
197
198PipeBuffer::Result PipeBuffer::transferTo(PipeBuffer& output, size_t maximum, bool consume,
199 bool block) {
200 if (this == &output) {
201 return {Status::Invalid, 0};
202 }
203 if (!maximum) {
204 return {Status::Ready, 0};
205 }
206 ActiveOperation inputOperation(*this);
207 ActiveOperation outputOperation(output);
208 if (!inputOperation) {
209 return {Status::Eof, 0};
210 }
211 if (!outputOperation) {
212 return {Status::Closed, 0};
213 }
214 PipeBuffer* first =
215 reinterpret_cast<uintptr_t>(this) < reinterpret_cast<uintptr_t>(&output) ? this : &output;
216 PipeBuffer* second = first == this ? &output : this;
217 PairWaiter waiter;
218 PairLink inputLink{&waiter};
219 PairLink outputLink{&waiter};
220
221 lock(first->m_Lock);
222 lock(second->m_Lock);
223 Result result{Status::WouldBlock, 0};
224 for (;;) {
225 const Result input = readyLocked(false);
226 const Result destination = output.readyLocked(true);
227 if (destination.status == Status::Closed || input.status == Status::Eof) {
228 result = destination.status == Status::Closed ? destination : input;
229 break;
230 }
231 if (input.status == Status::Ready && destination.status == Status::Ready) {
232 size_t accepted = maximum < input.count ? maximum : input.count;
233 accepted = accepted < destination.count ? accepted : destination.count;
234 const size_t tail = (output.m_Head + output.m_Size) % Capacity;
235 for (size_t i = 0; i < accepted; ++i) {
236 output.m_Data[(tail + i) % Capacity] = m_Data[(m_Head + i) % Capacity];
237 }
238 output.m_Size += accepted;
239 output.changedLocked();
240 if (consume) {
241 consumeLocked(accepted);
242 changedLocked();
243 }
244 result = {Status::Ready, accepted};
245 break;
246 }
247 if (!block) {
248 break;
249 }
250 if (interrupted()) {
251 result = {Status::Interrupted, 0};
252 break;
253 }
254 linkPairLocked(inputLink);
255 output.linkPairLocked(outputLink);
256 second->m_Lock.release();
257 first->m_Lock.release();
258
259 // The predicate and enrollment share the pair locks. Sticky notification
260 // bridges their release and this wait without an inverse waiter→buffer lock.
261 lock(waiter.mutex);
262 bool success = true;
263 while (!waiter.pending && success) {
264 success = !interrupted() && wait(waiter.condition, waiter.mutex);
265 }
266 waiter.pending = false;
267 waiter.mutex.release();
268 lock(first->m_Lock);
269 lock(second->m_Lock);
270 unlinkPairLocked(inputLink);
271 output.unlinkPairLocked(outputLink);
272 if (!success) {
273 result = {Status::Interrupted, 0};
274 break;
275 }
276 }
277 second->m_Lock.release();
278 first->m_Lock.release();
279 if (result.count) {
280 output.publish();
281 if (consume) {
282 publish();
283 }
284 }
285 return result;
286}
Definition Mutex.h:56
void release(size_t n=1)
Definition Semaphore.cc:549
#define assert(x)
Definition assert.h:39