3#include "pedigree/kernel/Log.h"
4#include "pedigree/kernel/utilities/assert.h"
5#include "pedigree/kernel/utilities/utility.h"
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"
13PipeBuffer::PipeBuffer(ChangeCallback changed,
void* context)
14 : m_ReadGeneration(0), m_WriteGeneration(0), m_Changed(changed), m_Context(context) {}
16PipeBuffer::~PipeBuffer() {
20 while (m_ActiveOperations) {
23 assert(!m_PairWaiters && !m_ReadReserved && !m_WriteReserved);
27void PipeBuffer::lock(
Mutex& mutex) {
28#ifdef STANDALONE_MUTEXES
29 const bool acquired = mutex.acquire();
31 const bool acquired = mutex.acquireForCompletion();
34 FATAL(
"PipeBuffer could not acquire its lifetime mutex.");
39 ConditionVariable::Error error = ConditionVariable::NoError;
40 const bool success = condition.
wait(mutex, error);
44#if THREADS && !defined(PEDIGREE_BUILDUTILS)
45 if (error == ConditionVariable::Interrupted) {
52bool PipeBuffer::interrupted() {
53#if THREADS && !defined(PEDIGREE_BUILDUTILS)
55 return thread && (thread->getInterruptionReason() == Thread::InterruptedBySignal ||
62PipeBuffer::ActiveOperation::ActiveOperation(
PipeBuffer& buffer)
63 : m_Termination(), m_Buffer(buffer.beginOperation() ? &buffer : nullptr) {}
65PipeBuffer::ActiveOperation::~ActiveOperation() {
67 m_Buffer->endOperation();
71PipeBuffer::ActiveOperation::operator bool()
const {
72 return m_Buffer !=
nullptr;
75void PipeBuffer::ActiveOperation::detach() {
79bool PipeBuffer::beginOperation() {
81 const bool admitted = !m_Closing;
89void PipeBuffer::endOperation() {
91 assert(m_ActiveOperations);
92 if (!--m_ActiveOperations) {
98bool PipeBuffer::readableLocked()
const {
99 return !m_Closing && m_ReadEnabled && !m_ResetPending && !m_ReadReserved && m_Size;
102bool PipeBuffer::writableLocked()
const {
103 return !m_Closing && m_ReadEnabled && m_WriteEnabled && !m_ResetPending && !m_WriteReserved &&
109 if (m_Closing || !m_ReadEnabled || !m_WriteEnabled) {
110 return {Status::Closed, 0};
112 if (writableLocked() && Capacity - m_Size >= minimum) {
113 return {Status::Ready, Capacity - m_Size};
116 if (m_Closing || !m_ReadEnabled || (!m_WriteEnabled && !m_Size)) {
117 return {Status::Eof, 0};
119 if (readableLocked()) {
120 return {Status::Ready, m_Size};
123 return {Status::WouldBlock, 0};
128 const Result result = readyLocked(writing, minimum);
129 if (result.status != Status::WouldBlock || !block) {
132 if (interrupted() || !wait(writing ? m_WriteCondition : m_ReadCondition, m_Lock)) {
133 return {Status::Interrupted, 0};
138void PipeBuffer::changedLocked() {
139 const bool readable = readableLocked();
140 const bool writable = writableLocked();
141 if (readable && !m_WasReadable) {
142 m_ReadGeneration += 1;
144 if (writable && !m_WasWritable) {
145 m_WriteGeneration += 1;
147 m_WasReadable = readable;
148 m_WasWritable = writable;
154void PipeBuffer::resetLocked() {
155 m_ResetPending =
true;
159void PipeBuffer::finishResetLocked() {
160 if (m_ResetPending && !m_ReadReserved && !m_WriteReserved) {
163 m_ResetPending =
false;
167void PipeBuffer::copyOutLocked(uint8_t* destination,
size_t count)
const {
169 const size_t first = count < Capacity - m_Head ? count : Capacity - m_Head;
170 pedigree_std::copy(destination, m_Data + m_Head, first);
172 pedigree_std::copy(destination + first, m_Data, count - first);
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);
182 pedigree_std::copy(m_Data, source + first, count - first);
187void PipeBuffer::consumeLocked(
size_t count) {
189 m_Head = (m_Head + count) % Capacity;
193void PipeBuffer::publish() {
195 m_Changed(m_Context);
199size_t PipeBuffer::read(uint8_t* destination,
size_t count,
bool block) {
200 ActiveOperation operation(*
this);
201 if (!operation || !count) {
205 const Result result = waitLocked(
false, 1, block);
207 if (result.status == Status::Ready) {
208 accepted = count < result.count ? count : result.count;
209 copyOutLocked(destination, accepted);
210 consumeLocked(accepted);
220size_t PipeBuffer::write(
const uint8_t* source,
size_t count,
bool block) {
221 ActiveOperation operation(*
this);
226 while (written < count) {
231 const Result result = waitLocked(
true, 1, block);
232 if (result.status != Status::Ready) {
236 const size_t remaining = count - written;
237 const size_t accepted = remaining < result.count ? remaining : result.count;
238 appendLocked(source + written, accepted);
248size_t PipeBuffer::writeAtomic(
const uint8_t* source,
size_t count,
bool block) {
249 ActiveOperation operation(*
this);
250 if (!operation || !count || count > Capacity) {
254 const Result result = waitLocked(
true, count, block);
255 if (result.status == Status::Ready) {
256 appendLocked(source, count);
260 if (result.status != Status::Ready) {
268 ActiveOperation operation(*
this);
270 return {writing ? Status::Closed : Status::Eof, 0};
273 const Result result = waitLocked(writing, 1, block);
278bool PipeBuffer::canRead(
bool block) {
279 return waitTransfer(
false, block).status == Status::Ready;
282bool PipeBuffer::canWrite(
bool block) {
283 return waitTransfer(
true, block).status == Status::Ready;
286uint64_t PipeBuffer::readableGeneration()
const {
287 return m_ReadGeneration.value();
290uint64_t PipeBuffer::writableGeneration()
const {
291 return m_WriteGeneration.value();
294size_t PipeBuffer::getDataSize() {
297 const size_t size = m_Size;
302void PipeBuffer::disableReads() {
303 ActiveOperation operation(*
this);
306 m_ReadEnabled =
false;
312void PipeBuffer::disableWrites() {
313 ActiveOperation operation(*
this);
316 m_WriteEnabled =
false;
322bool PipeBuffer::enableReads() {
323 ActiveOperation operation(*
this);
328 const bool previous = m_ReadEnabled;
329 m_ReadEnabled =
true;
335bool PipeBuffer::enableWrites(
bool resetIfPreviouslyDisabled) {
336 ActiveOperation operation(*
this);
341 const bool previous = m_WriteEnabled;
342 m_WriteEnabled =
true;
343 if (!previous && resetIfPreviouslyDisabled) {
351void PipeBuffer::wipe() {
352 ActiveOperation operation(*
this);
361void PipeBuffer::close() {
365 m_ReadEnabled =
false;
366 m_WriteEnabled =
false;
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)
static ProcessorInformation & information()
@ Continue
No unwind necessary, carry on as normal.
UnwindType getUnwindState()