The Pedigree Project 0.1
PipeBuffer.cc
1/* Copyright (c) 2026, Pedigree Developers. */
2#include "PipeBuffer.h"
3#include "pedigree/kernel/Log.h"
4#include "pedigree/kernel/utilities/assert.h"
5#include "pedigree/kernel/utilities/utility.h"
6
7#if THREADS && !defined(PEDIGREE_BUILDUTILS)
8#include "pedigree/kernel/process/Thread.h"
9#include "pedigree/kernel/processor/Processor.h"
10#include "pedigree/kernel/processor/ProcessorInformation.h"
11#endif
12
13PipeBuffer::PipeBuffer(ChangeCallback changed, void* context)
14 : m_ReadGeneration(0), m_WriteGeneration(0), m_Changed(changed), m_Context(context) {}
15
16PipeBuffer::~PipeBuffer() {
17 TerminationDeferral termination;
18 close();
19 lock(m_Lock);
20 while (m_ActiveOperations) {
21 m_DrainCondition.waitForCompletion(m_Lock);
22 }
23 assert(!m_PairWaiters && !m_ReadReserved && !m_WriteReserved);
24 m_Lock.release();
25}
26
27void PipeBuffer::lock(Mutex& mutex) {
28#ifdef STANDALONE_MUTEXES
29 const bool acquired = mutex.acquire();
30#else
31 const bool acquired = mutex.acquireForCompletion();
32#endif
33 if (!acquired) {
34 FATAL("PipeBuffer could not acquire its lifetime mutex.");
35 }
36}
37
38bool PipeBuffer::wait(ConditionVariable& condition, Mutex& mutex) {
39 ConditionVariable::Error error = ConditionVariable::NoError;
40 const bool success = condition.wait(mutex, error);
42 lock(mutex);
43 }
44#if THREADS && !defined(PEDIGREE_BUILDUTILS)
45 if (error == ConditionVariable::Interrupted) {
46 Processor::information().getCurrentThread()->setInterruptionReason(Thread::InterruptedBySignal);
47 }
48#endif
49 return success;
50}
51
52bool PipeBuffer::interrupted() {
53#if THREADS && !defined(PEDIGREE_BUILDUTILS)
54 Thread* thread = Processor::information().getCurrentThread();
55 return thread && (thread->getInterruptionReason() == Thread::InterruptedBySignal ||
56 thread->getUnwindState() != Thread::Continue);
57#else
58 return false;
59#endif
60}
61
62PipeBuffer::ActiveOperation::ActiveOperation(PipeBuffer& buffer)
63 : m_Termination(), m_Buffer(buffer.beginOperation() ? &buffer : nullptr) {}
64
65PipeBuffer::ActiveOperation::~ActiveOperation() {
66 if (m_Buffer) {
67 m_Buffer->endOperation();
68 }
69}
70
71PipeBuffer::ActiveOperation::operator bool() const {
72 return m_Buffer != nullptr;
73}
74
75void PipeBuffer::ActiveOperation::detach() {
76 m_Buffer = nullptr;
77}
78
79bool PipeBuffer::beginOperation() {
80 lock(m_Lock);
81 const bool admitted = !m_Closing;
82 if (admitted) {
83 ++m_ActiveOperations;
84 }
85 m_Lock.release();
86 return admitted;
87}
88
89void PipeBuffer::endOperation() {
90 lock(m_Lock);
91 assert(m_ActiveOperations);
92 if (!--m_ActiveOperations) {
93 m_DrainCondition.broadcast();
94 }
95 m_Lock.release();
96}
97
98bool PipeBuffer::readableLocked() const {
99 return !m_Closing && m_ReadEnabled && !m_ResetPending && !m_ReadReserved && m_Size;
100}
101
102bool PipeBuffer::writableLocked() const {
103 return !m_Closing && m_ReadEnabled && m_WriteEnabled && !m_ResetPending && !m_WriteReserved &&
104 m_Size < Capacity;
105}
106
107PipeBuffer::Result PipeBuffer::readyLocked(bool writing, size_t minimum) const {
108 if (writing) {
109 if (m_Closing || !m_ReadEnabled || !m_WriteEnabled) {
110 return {Status::Closed, 0};
111 }
112 if (writableLocked() && Capacity - m_Size >= minimum) {
113 return {Status::Ready, Capacity - m_Size};
114 }
115 } else {
116 if (m_Closing || !m_ReadEnabled || (!m_WriteEnabled && !m_Size)) {
117 return {Status::Eof, 0};
118 }
119 if (readableLocked()) {
120 return {Status::Ready, m_Size};
121 }
122 }
123 return {Status::WouldBlock, 0};
124}
125
126PipeBuffer::Result PipeBuffer::waitLocked(bool writing, size_t minimum, bool block) {
127 for (;;) {
128 const Result result = readyLocked(writing, minimum);
129 if (result.status != Status::WouldBlock || !block) {
130 return result;
131 }
132 if (interrupted() || !wait(writing ? m_WriteCondition : m_ReadCondition, m_Lock)) {
133 return {Status::Interrupted, 0};
134 }
135 }
136}
137
138void PipeBuffer::changedLocked() {
139 const bool readable = readableLocked();
140 const bool writable = writableLocked();
141 if (readable && !m_WasReadable) {
142 m_ReadGeneration += 1;
143 }
144 if (writable && !m_WasWritable) {
145 m_WriteGeneration += 1;
146 }
147 m_WasReadable = readable;
148 m_WasWritable = writable;
149 m_ReadCondition.broadcast();
150 m_WriteCondition.broadcast();
151 wakePairsLocked();
152}
153
154void PipeBuffer::resetLocked() {
155 m_ResetPending = true;
156 finishResetLocked();
157}
158
159void PipeBuffer::finishResetLocked() {
160 if (m_ResetPending && !m_ReadReserved && !m_WriteReserved) {
161 m_Head = 0;
162 m_Size = 0;
163 m_ResetPending = false;
164 }
165}
166
167void PipeBuffer::copyOutLocked(uint8_t* destination, size_t count) const {
168 assert(count <= m_Size);
169 const size_t first = count < Capacity - m_Head ? count : Capacity - m_Head;
170 pedigree_std::copy(destination, m_Data + m_Head, first);
171 if (count > first) {
172 pedigree_std::copy(destination + first, m_Data, count - first);
173 }
174}
175
176void PipeBuffer::appendLocked(const uint8_t* source, size_t count) {
177 assert(count <= Capacity - m_Size);
178 const size_t tail = (m_Head + m_Size) % Capacity;
179 const size_t first = count < Capacity - tail ? count : Capacity - tail;
180 pedigree_std::copy(m_Data + tail, source, first);
181 if (count > first) {
182 pedigree_std::copy(m_Data, source + first, count - first);
183 }
184 m_Size += count;
185}
186
187void PipeBuffer::consumeLocked(size_t count) {
188 assert(count <= m_Size);
189 m_Head = (m_Head + count) % Capacity;
190 m_Size -= count;
191}
192
193void PipeBuffer::publish() {
194 if (m_Changed) {
195 m_Changed(m_Context);
196 }
197}
198
199size_t PipeBuffer::read(uint8_t* destination, size_t count, bool block) {
200 ActiveOperation operation(*this);
201 if (!operation || !count) {
202 return 0;
203 }
204 lock(m_Lock);
205 const Result result = waitLocked(false, 1, block);
206 size_t accepted = 0;
207 if (result.status == Status::Ready) {
208 accepted = count < result.count ? count : result.count;
209 copyOutLocked(destination, accepted);
210 consumeLocked(accepted);
211 changedLocked();
212 }
213 m_Lock.release();
214 if (accepted) {
215 publish();
216 }
217 return accepted;
218}
219
220size_t PipeBuffer::write(const uint8_t* source, size_t count, bool block) {
221 ActiveOperation operation(*this);
222 if (!operation) {
223 return 0;
224 }
225 size_t written = 0;
226 while (written < count) {
227 if (interrupted()) {
228 break;
229 }
230 lock(m_Lock);
231 const Result result = waitLocked(true, 1, block);
232 if (result.status != Status::Ready) {
233 m_Lock.release();
234 break;
235 }
236 const size_t remaining = count - written;
237 const size_t accepted = remaining < result.count ? remaining : result.count;
238 appendLocked(source + written, accepted);
239 changedLocked();
240 m_Lock.release();
241 written += accepted;
242 // A readiness-driven reader must see this chunk before we wait for space.
243 publish();
244 }
245 return written;
246}
247
248size_t PipeBuffer::writeAtomic(const uint8_t* source, size_t count, bool block) {
249 ActiveOperation operation(*this);
250 if (!operation || !count || count > Capacity) {
251 return 0;
252 }
253 lock(m_Lock);
254 const Result result = waitLocked(true, count, block);
255 if (result.status == Status::Ready) {
256 appendLocked(source, count);
257 changedLocked();
258 }
259 m_Lock.release();
260 if (result.status != Status::Ready) {
261 return 0;
262 }
263 publish();
264 return count;
265}
266
267PipeBuffer::Result PipeBuffer::waitTransfer(bool writing, bool block) {
268 ActiveOperation operation(*this);
269 if (!operation) {
270 return {writing ? Status::Closed : Status::Eof, 0};
271 }
272 lock(m_Lock);
273 const Result result = waitLocked(writing, 1, block);
274 m_Lock.release();
275 return result;
276}
277
278bool PipeBuffer::canRead(bool block) {
279 return waitTransfer(false, block).status == Status::Ready;
280}
281
282bool PipeBuffer::canWrite(bool block) {
283 return waitTransfer(true, block).status == Status::Ready;
284}
285
286uint64_t PipeBuffer::readableGeneration() const {
287 return m_ReadGeneration.value();
288}
289
290uint64_t PipeBuffer::writableGeneration() const {
291 return m_WriteGeneration.value();
292}
293
294size_t PipeBuffer::getDataSize() {
295 TerminationDeferral termination;
296 lock(m_Lock);
297 const size_t size = m_Size;
298 m_Lock.release();
299 return size;
300}
301
302void PipeBuffer::disableReads() {
303 ActiveOperation operation(*this);
304 if (operation) {
305 lock(m_Lock);
306 m_ReadEnabled = false;
307 changedLocked();
308 m_Lock.release();
309 }
310}
311
312void PipeBuffer::disableWrites() {
313 ActiveOperation operation(*this);
314 if (operation) {
315 lock(m_Lock);
316 m_WriteEnabled = false;
317 changedLocked();
318 m_Lock.release();
319 }
320}
321
322bool PipeBuffer::enableReads() {
323 ActiveOperation operation(*this);
324 if (!operation) {
325 return false;
326 }
327 lock(m_Lock);
328 const bool previous = m_ReadEnabled;
329 m_ReadEnabled = true;
330 changedLocked();
331 m_Lock.release();
332 return previous;
333}
334
335bool PipeBuffer::enableWrites(bool resetIfPreviouslyDisabled) {
336 ActiveOperation operation(*this);
337 if (!operation) {
338 return false;
339 }
340 lock(m_Lock);
341 const bool previous = m_WriteEnabled;
342 m_WriteEnabled = true;
343 if (!previous && resetIfPreviouslyDisabled) {
344 resetLocked();
345 }
346 changedLocked();
347 m_Lock.release();
348 return previous;
349}
350
351void PipeBuffer::wipe() {
352 ActiveOperation operation(*this);
353 if (operation) {
354 lock(m_Lock);
355 resetLocked();
356 changedLocked();
357 m_Lock.release();
358 }
359}
360
361void PipeBuffer::close() {
362 TerminationDeferral termination;
363 lock(m_Lock);
364 m_Closing = true;
365 m_ReadEnabled = false;
366 m_WriteEnabled = false;
367 changedLocked();
368 m_Lock.release();
369}
void waitForCompletion(Mutex &mutex)
MUST_USE_RESULT bool wait(Mutex &mutex, Time::Timestamp &timeout, Error &error, WaitQueue::StackDiscardCleanup onStackDiscard=nullptr, void *stackDiscardContext=nullptr)
static bool mutexAcquired(Error error)
Definition Mutex.h:56
static ProcessorInformation & information()
void release(size_t n=1)
Definition Semaphore.cc:549
@ Continue
No unwind necessary, carry on as normal.
Definition Thread.h:500
UnwindType getUnwindState()
Definition Thread.h:518
#define assert(x)
Definition assert.h:39