2#include "pedigree/kernel/utilities/assert.h"
6PipeBuffer::ReadReservation::ReadReservation() =
default;
8PipeBuffer::ReadReservation::~ReadReservation() {
12size_t PipeBuffer::ReadReservation::size()
const {
16void PipeBuffer::ReadReservation::copyTo(uint8_t* destination,
size_t count)
const {
22 PipeBuffer::lock(m_Buffer->m_Lock);
23 m_Buffer->copyOutLocked(destination, count);
24 m_Buffer->m_Lock.release();
27void PipeBuffer::ReadReservation::consume(
size_t accepted) {
28 assert(accepted <= m_Size);
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();
43 buffer->endOperation();
46void PipeBuffer::ReadReservation::cancel() {
50PipeBuffer::WriteReservation::WriteReservation() =
default;
52PipeBuffer::WriteReservation::~WriteReservation() {
56size_t PipeBuffer::WriteReservation::size()
const {
60PipeBuffer::Result PipeBuffer::WriteReservation::commit(
const uint8_t* source,
size_t accepted) {
63 return {Status::Invalid, 0};
65 if (accepted > m_Size) {
67 return {Status::Invalid, 0};
69 PipeBuffer::lock(buffer->m_Lock);
70 assert(buffer->m_WriteReserved);
72 const bool closed = buffer->m_Closing || !buffer->m_ReadEnabled || !buffer->m_WriteEnabled ||
73 buffer->m_ResetPending;
75 buffer->appendLocked(source, accepted);
77 buffer->m_WriteReserved =
false;
78 buffer->finishResetLocked();
79 buffer->changedLocked();
84 buffer->endOperation();
85 return {closed ? Status::Closed : Status::Ready, closed ? 0 : accepted};
88void PipeBuffer::WriteReservation::cancel() {
93 PipeBuffer::lock(buffer->m_Lock);
94 assert(buffer->m_WriteReserved);
95 buffer->m_WriteReserved =
false;
96 buffer->finishResetLocked();
97 buffer->changedLocked();
102 buffer->endOperation();
106 ReadReservation& reservation) {
107 if (reservation.m_Buffer) {
108 return {Status::Invalid, 0};
111 return {Status::Ready, 0};
113 ActiveOperation operation(*
this);
115 return {Status::Eof, 0};
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;
128 if (result.status == Status::Ready) {
135 WriteReservation& reservation) {
136 if (reservation.m_Buffer) {
137 return {Status::Invalid, 0};
140 return {Status::Ready, 0};
142 ActiveOperation operation(*
this);
144 return {Status::Closed, 0};
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;
157 if (result.status == Status::Ready) {
166 bool pending =
false;
174void PipeBuffer::linkPairLocked(
PairLink& link) {
175 link.next = m_PairWaiters;
176 m_PairWaiters = &link;
179void PipeBuffer::unlinkPairLocked(PairLink& link) {
180 PairLink** position = &m_PairWaiters;
181 while (*position && *position != &link) {
182 position = &(*position)->next;
184 assert(*position == &link);
185 *position = link.next;
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();
200 if (
this == &output) {
201 return {Status::Invalid, 0};
204 return {Status::Ready, 0};
206 ActiveOperation inputOperation(*
this);
207 ActiveOperation outputOperation(output);
208 if (!inputOperation) {
209 return {Status::Eof, 0};
211 if (!outputOperation) {
212 return {Status::Closed, 0};
215 reinterpret_cast<uintptr_t
>(
this) <
reinterpret_cast<uintptr_t
>(&output) ? this : &output;
216 PipeBuffer* second = first ==
this ? &output :
this;
218 PairLink inputLink{&
waiter};
219 PairLink outputLink{&
waiter};
222 lock(second->m_Lock);
223 Result result{Status::WouldBlock, 0};
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;
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];
238 output.m_Size += accepted;
239 output.changedLocked();
241 consumeLocked(accepted);
244 result = {Status::Ready, accepted};
251 result = {Status::Interrupted, 0};
254 linkPairLocked(inputLink);
255 output.linkPairLocked(outputLink);
263 while (!
waiter.pending && success) {
264 success = !interrupted() && wait(
waiter.condition,
waiter.mutex);
269 lock(second->m_Lock);
270 unlinkPairLocked(inputLink);
271 output.unlinkPairLocked(outputLink);
273 result = {Status::Interrupted, 0};