The Pedigree Project 0.1
Buffer.cc
1/*
2 * Copyright (c) 2008-2014, Pedigree Developers
3 *
4 * Please see the CONTRIB file in the root of the source tree for a full
5 * list of contributors.
6 *
7 * Permission to use, copy, modify, and distribute this software for any
8 * purpose with or without fee is hereby granted, provided that the above
9 * copyright notice and this permission notice appear in all copies.
10 *
11 * THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
12 * WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
13 * MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
14 * ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
15 * WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
16 * ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
17 * OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
18 */
19
20#include "pedigree/kernel/LockGuard.h"
21#include "pedigree/kernel/Log.h"
22#include "pedigree/kernel/utilities/Buffer.h"
23#include "pedigree/kernel/utilities/assert.h"
24#include "pedigree/kernel/utilities/utility.h"
25
26#if THREADS
27#include "pedigree/kernel/process/Thread.h"
28#include "pedigree/kernel/processor/Processor.h"
29#include "pedigree/kernel/processor/ProcessorInformation.h"
30
31namespace {
32void preserveConditionInterruption(ConditionVariable::Error error) {
33 if (error == ConditionVariable::Interrupted) {
34 Processor::information().getCurrentThread()->setInterruptionReason(Thread::InterruptedBySignal);
35 }
36}
37} // namespace
38#endif
39
40template <class T, bool allowShortOperation>
42 : m_BufferSize(bufferSize),
43 m_DataSize(0),
44 m_ReadableGeneration(0),
45 m_WritableGeneration(0),
46 m_Lock(),
47 m_WriteCondition(),
48 m_ReadCondition(),
49 m_DrainCondition(),
50 m_Segments(),
51 m_Monitors(),
52 m_bCanRead(true),
53 m_bCanWrite(true),
54 m_bClosing(false),
55 m_ActiveOperations(0) {}
56
57template <class T, bool allowShortOperation>
59 TerminationDeferral terminationDeferral;
60 m_Lock.acquire();
61 m_bClosing = true;
62 m_bCanRead = false;
63 m_bCanWrite = false;
64 m_ReadCondition.broadcast();
65 m_WriteCondition.broadcast();
66
67 while (m_ActiveOperations) {
68 m_DrainCondition.waitForCompletion(m_Lock);
69 }
70
71 for (auto pSegment : m_Segments) {
72 delete pSegment;
73 }
74 m_Segments.clear();
75 m_DataSize = 0;
76
77 m_Monitors.clear();
78 m_Lock.release();
79}
80
81template <class T, bool allowShortOperation>
83 LockGuard<Mutex> guard(m_Lock);
84 if (m_bClosing) {
85 return false;
86 }
87
88 ++m_ActiveOperations;
89 return true;
90}
92template <class T, bool allowShortOperation>
94 LockGuard<Mutex> guard(m_Lock);
95 assert(m_ActiveOperations);
96 --m_ActiveOperations;
97 if (m_bClosing && !m_ActiveOperations) {
98 m_DrainCondition.signal();
99 }
100}
101
102template <class T, bool allowShortOperation>
103size_t Buffer<T, allowShortOperation>::write(const T* buffer, size_t count, bool block) {
104 ActiveOperation operation(*this);
105 if (!operation) {
106 return 0;
107 }
108
109 if (!block) {
110 if (!m_Lock.tryAcquire()) {
111 // can't unlock buffer for writing
112 return 0;
113 }
114 } else {
115 // can block!
116 m_Lock.acquire();
118
119 size_t countSoFar = 0;
120 while (true) {
121 // Can we write?
122 if (!m_bCanWrite) {
123 // No! Maybe not anymore, so return what we've written so far.
124 break;
125 }
126
127 // Do we have space?
128 size_t bytesAvailable = m_BufferSize - m_DataSize;
129 if (!bytesAvailable) {
130 if (!block) {
131 // Cannot block!
132 break;
134
135 // Can any reader get us out of this situation?
136 if (!m_bCanRead) {
137 // No. Return what we've written so far.
138 break;
139 }
140
141 // No, we need to wait.
142 ConditionVariable::Error error = ConditionVariable::NoError;
143 if (!m_WriteCondition.wait(m_Lock, error)) {
144#if THREADS
145 preserveConditionInterruption(error);
146#endif
148 m_Lock.release();
149 }
150 return countSoFar;
151 }
152 continue;
153 }
155 // Yes, we have room.
156 size_t totalCount = bytesAvailable;
157 if (totalCount > count) {
158 totalCount = count;
160
161 // If we're allowed to just give up if we don't have enough bytes, do
162 // so (e.g. for TCP buffers and the like).
163 if (allowShortOperation && (totalCount > bytesAvailable)) {
164 totalCount = bytesAvailable;
165 count = bytesAvailable;
166 if (!totalCount) {
167 break;
168 }
170
171 const size_t numberCopied = writeLocked(buffer, totalCount);
172 countSoFar += numberCopied;
173 buffer += numberCopied;
174 count -= numberCopied;
175
176 // Wake up a reader that was waiting for data.
177 // We do this here rather than after we finish as we may need to block
178 // again (e.g. if we're writing more than the size of the buffer), and
179 // a reader is needed to unblock that.
180 if (numberCopied) {
181 m_ReadCondition.signal();
182 }
183
184 if (!count) {
185 // Complete.
186 break;
187 }
188 }
189
190 m_Lock.release();
191
192 if (countSoFar) {
193 // We've updated the buffer, so send events.
194 notifyMonitors();
195 }
196
197 return countSoFar;
198}
199
200template <class T, bool allowShortOperation>
201size_t Buffer<T, allowShortOperation>::writeAvailable(const T* buffer, size_t count, bool atomic) {
202 ActiveOperation operation(*this);
203 if (!operation) {
204 return 0;
205 }
206
207 m_Lock.acquire();
208 size_t written = 0;
209 if (m_bCanRead && m_bCanWrite) {
210 const size_t available = m_BufferSize - m_DataSize;
211 if (!atomic || count <= available) {
212 written = writeLocked(buffer, count < available ? count : available);
213 if (written) {
214 m_ReadCondition.signal();
215 }
216 }
218 m_Lock.release();
219
220 if (written) {
221 notifyMonitors();
222 }
223 return written;
224}
225
226template <class T, bool allowShortOperation>
227size_t Buffer<T, allowShortOperation>::writeAtomic(const T* buffer, size_t count, bool block) {
228 ActiveOperation operation(*this);
229 if (!operation || count > m_BufferSize) {
230 return 0;
231 }
232
233 if (!block) {
234 if (!m_Lock.tryAcquire()) {
235 return 0;
236 }
237 } else {
238 m_Lock.acquire();
239 }
240
241 size_t written = 0;
242 while (m_bCanWrite && m_bCanRead) {
243 if (count <= (m_BufferSize - m_DataSize)) {
244 written = writeLocked(buffer, count);
245 if (written) {
246 m_ReadCondition.signal();
247 }
248 break;
249 }
250
251 if (!block || !m_bCanRead) {
252 break;
253 }
254
255 ConditionVariable::Error error = ConditionVariable::NoError;
256 if (!m_WriteCondition.wait(m_Lock, error)) {
257#if THREADS
258 preserveConditionInterruption(error);
259#endif
261 m_Lock.release();
262 }
263 return 0;
264 }
265 }
266
267 m_Lock.release();
268 if (written) {
269 notifyMonitors();
270 }
271 return written;
272}
273
274template <class T, bool allowShortOperation>
275bool Buffer<T, allowShortOperation>::tryWrite(const T* buffer, size_t count) {
276 TerminationDeferral terminationDeferral;
277 if (!m_Lock.tryAcquire()) {
278 return false;
279 }
280
281 if (m_bClosing || !m_bCanWrite || count > (m_BufferSize - m_DataSize)) {
282 m_Lock.release();
283 return false;
284 }
285
286 const size_t numberCopied = writeLocked(buffer, count);
287 if (numberCopied) {
288 m_ReadCondition.signal();
289 notifyMonitorsLocked();
290 }
291 m_Lock.release();
292 return numberCopied == count;
293}
294
295template <class T, bool allowShortOperation>
296size_t Buffer<T, allowShortOperation>::read(T* buffer, size_t count, bool block) {
297 ActiveOperation operation(*this);
298 if (!operation) {
299 return 0;
300 }
301
302 // Nonblocking reads must still wait for lock ownership, as beginOperation
303 // already does. Contention is not an empty buffer (or EOF to the caller).
304 m_Lock.acquire();
305
306 size_t countSoFar = 0;
307 while (true) {
308 // Can we read?
309 if (!m_bCanRead) {
310 // No! Maybe not anymore, so return what we've read so far.
311 break;
312 }
313
314 // Do we have anything to read?
315 if (!m_DataSize) {
316 if (!block) {
317 // Cannot block!
318 break;
319 }
320
321 // Can any writer get us out of this situation?
322 if (!m_bCanWrite) {
323 // No. Return what we've read so far.
324 break;
325 }
326
327 // No, we need to wait.
328 ConditionVariable::Error error = ConditionVariable::NoError;
329 if (!m_ReadCondition.wait(m_Lock, error)) {
330#if THREADS
331 preserveConditionInterruption(error);
332#endif
334 m_Lock.release();
335 }
336 return countSoFar;
337 }
338 continue;
339 }
340
341 // Yes, we have room.
342 size_t totalCount = count;
343 if (totalCount > m_DataSize) {
344 totalCount = m_DataSize;
345 }
346
347 size_t numberCopied = 0;
348 while (m_Segments.count() && numberCopied < totalCount) {
349 // Grab the first segment and read it.
350 Segment* pSegment = m_Segments.popFront();
351 size_t countToRead = pSegment->size - pSegment->reader;
352 if ((numberCopied + countToRead) > totalCount) {
353 countToRead = totalCount - numberCopied;
354 }
355
356 // Copy.
357 pedigree_std::copy(buffer, &pSegment->data[pSegment->reader], countToRead);
358 pSegment->reader += countToRead;
359
360 // Do we need to re-add it?
361 if (pSegment->reader < pSegment->size) {
362 m_Segments.pushFront(pSegment);
363 } else {
364 delete pSegment;
365 }
366
367 numberCopied += countToRead;
368 buffer += countToRead;
369 }
370
371 const bool wasFull = m_DataSize >= m_BufferSize;
372 m_DataSize -= numberCopied;
373 if (wasFull && m_DataSize < m_BufferSize) {
374 m_WritableGeneration += 1;
375 }
376 countSoFar += numberCopied;
377 count -= numberCopied;
378
379 // We read some bytes so writers may be able to continue, which may be
380 // needed to unblock us if we loop back around and block.
381 if (numberCopied) {
382 // Wake up a writer that was waiting for space to write.
383 m_WriteCondition.signal();
384 }
385
386 if (!count) {
387 break;
388 }
389
390 // Once we've read at least some bytes, don't block - just return what
391 // we've read so far if we loop back around and have no data.
392 block = false;
393 }
394
395 m_Lock.release();
396
397 if (countSoFar) {
398 // We've updated the buffer, so send events.
399 notifyMonitors();
400 }
401
402 return countSoFar;
403}
404
405template <class T, bool allowShortOperation>
407 ActiveOperation operation(*this);
408 if (!operation) {
409 return;
410 }
411
412 LockGuard<Mutex> guard(m_Lock);
413 m_bCanWrite = false;
414
415 // Readers may be waiting for data, while writers may already be waiting
416 // for space. Both predicates have changed permanently for this direction.
417 m_ReadCondition.broadcast();
418 m_WriteCondition.broadcast();
419}
420
421template <class T, bool allowShortOperation>
423 ActiveOperation operation(*this);
424 if (!operation) {
425 return;
426 }
427
428 LockGuard<Mutex> guard(m_Lock);
429 m_bCanRead = false;
430
431 // Writers may be waiting for space, while readers may already be waiting
432 // for data. Both predicates have changed permanently for this direction.
433 m_ReadCondition.broadcast();
434 m_WriteCondition.broadcast();
435}
436
437template <class T, bool allowShortOperation>
439 ActiveOperation operation(*this);
440 if (!operation) {
441 return false;
442 }
443
444 LockGuard<Mutex> guard(m_Lock);
445 bool previous = m_bCanWrite;
446 m_bCanWrite = true;
447 if (!previous && m_DataSize < m_BufferSize) {
448 m_WritableGeneration += 1;
449 }
450 return previous;
451}
452
453template <class T, bool allowShortOperation>
455 ActiveOperation operation(*this);
456 if (!operation) {
457 return false;
458 }
459
460 LockGuard<Mutex> guard(m_Lock);
461 bool previous = m_bCanRead;
462 m_bCanRead = true;
463 if (!previous && m_DataSize) {
464 m_ReadableGeneration += 1;
465 }
466 return previous;
467}
468
469template <class T, bool allowShortOperation>
471 ActiveOperation operation(*this);
472 if (!operation) {
473 return 0;
474 }
475
476 LockGuard<Mutex> guard(m_Lock);
477 return m_DataSize;
478}
479
480template <class T, bool allowShortOperation>
482 ActiveOperation operation(*this);
483 if (!operation) {
484 return 0;
485 }
486
487 return m_BufferSize;
488}
489
490template <class T, bool allowShortOperation>
492 ActiveOperation operation(*this);
493 if (!operation) {
494 return false;
495 }
496
497 LockGuard<Mutex> guard(m_Lock);
498
499 if (!block) {
500 return m_bCanWrite && (m_DataSize < m_BufferSize);
501 }
502
503 if (!m_bCanWrite) {
504 return false;
505 }
506
507 // A full buffer can only become writable while reads remain possible.
508 while (m_bCanWrite && m_bCanRead && m_DataSize >= m_BufferSize) {
509 ConditionVariable::Error error = ConditionVariable::NoError;
510 if (!m_WriteCondition.wait(m_Lock, error)) {
511#if THREADS
512 preserveConditionInterruption(error);
513#endif
515 guard.disown();
516 }
517 return false;
518 }
519 }
520
521 return m_bCanWrite && m_DataSize < m_BufferSize;
522}
523
524template <class T, bool allowShortOperation>
526 ActiveOperation operation(*this);
527 if (!operation) {
528 return false;
529 }
530
531 LockGuard<Mutex> guard(m_Lock);
532
533 if (!block) {
534 return m_bCanRead && (m_DataSize > 0);
535 }
536
537 if (!m_bCanRead) {
538 return false;
539 }
540
541 // An empty buffer can only become readable while writes remain possible.
542 while (m_bCanRead && m_bCanWrite && !m_DataSize) {
543 ConditionVariable::Error error = ConditionVariable::NoError;
544 if (!m_ReadCondition.wait(m_Lock, error)) {
545#if THREADS
546 preserveConditionInterruption(error);
547#endif
549 guard.disown();
550 }
551 return false;
552 }
553 }
554
555 return m_bCanRead && m_DataSize > 0;
556}
557
558template <class T, bool allowShortOperation>
560 return m_ReadableGeneration.value();
561}
562
563template <class T, bool allowShortOperation>
565 return m_WritableGeneration.value();
566}
567
568template <class T, bool allowShortOperation>
570 ActiveOperation operation(*this);
571 if (!operation) {
572 return;
573 }
574
575 LockGuard<Mutex> guard(m_Lock);
576
577 const bool wasFull = m_DataSize >= m_BufferSize;
578
579 // Wipe out every segment we own.
580 for (auto pSegment : m_Segments) {
581 delete pSegment;
582 }
583 m_Segments.clear();
584 m_DataSize = 0;
585 if (wasFull && m_BufferSize) {
586 m_WritableGeneration += 1;
587 }
588
589 // Notify writers that might have been waiting for space.
590 m_WriteCondition.signal();
591}
592
593template <class T, bool allowShortOperation>
595 ActiveOperation operation(*this);
596 if (!operation) {
597 return;
598 }
599
600#if THREADS
601 LockGuard<Mutex> guard(m_Lock);
602 m_Monitors.add(pThread, pEvent);
603#endif
604}
605
606template <class T, bool allowShortOperation>
608 ActiveOperation operation(*this);
609 if (!operation) {
610 return;
611 }
612
613#if THREADS
614 LockGuard<Mutex> guard(m_Lock);
615 m_Monitors.add(pSemaphore);
616#endif
617}
618
619template <class T, bool allowShortOperation>
621 ActiveOperation operation(*this);
622 if (!operation) {
623 return;
624 }
625
626#if THREADS
627 LockGuard<Mutex> guard(m_Lock);
628 m_Monitors.cull(pThread);
629#endif
630}
631
632template <class T, bool allowShortOperation>
634 ActiveOperation operation(*this);
635 if (!operation) {
636 return;
637 }
638
639#if THREADS
640 LockGuard<Mutex> guard(m_Lock);
641 m_Monitors.cull(pSemaphore);
642#endif
643}
644
645template <class T, bool allowShortOperation>
647 ActiveOperation operation(*this);
648 if (!operation) {
649 return;
650 }
651
652#if THREADS
653 LockGuard<Mutex> guard(m_Lock);
654 m_Monitors.cull(pEvent);
655#endif
656}
657
658template <class T, bool allowShortOperation>
660 ActiveOperation operation(*this);
661 if (!operation) {
662 return;
663 }
664
665#if THREADS
666 LockGuard<Mutex> guard(m_Lock);
667 notifyMonitorsLocked();
668#endif
669}
670
671template <class T, bool allowShortOperation>
673#if THREADS
674 m_Monitors.notify();
675#endif
676}
677
678template <class T, bool allowShortOperation>
679size_t Buffer<T, allowShortOperation>::writeLocked(const T* buffer, size_t count) {
680 const bool wasReadable = m_bCanRead && m_DataSize;
681 size_t countSoFar = 0;
682 while (count) {
683 size_t numberCopied = 0;
684 if (count >= m_SegmentSize) {
685 addSegment(buffer, m_SegmentSize);
686 numberCopied = m_SegmentSize;
687 } else if (!m_Segments.count()) {
688 addSegment(buffer, count);
689 numberCopied = count;
690 } else {
691 Segment* pSegment = m_Segments.popBack();
692 if (pSegment->size == m_SegmentSize) {
693 m_Segments.pushBack(pSegment);
694 addSegment(buffer, count);
695 numberCopied = count;
696 } else {
697 T* start = &pSegment->data[pSegment->size];
698 size_t availableSpace = m_SegmentSize - pSegment->size;
699 if (availableSpace > count) {
700 availableSpace = count;
701 }
702 pedigree_std::copy(start, buffer, availableSpace);
703 pSegment->size += availableSpace;
704 m_Segments.pushBack(pSegment);
705
706 if (availableSpace < count) {
707 addSegment(&buffer[availableSpace], count - availableSpace);
708 }
709 numberCopied = count;
710 }
711 }
712
713 countSoFar += numberCopied;
714 m_DataSize += numberCopied;
715 buffer += numberCopied;
716 count -= numberCopied;
717 }
718
719 if (!wasReadable && m_bCanRead && m_DataSize) {
720 m_ReadableGeneration += 1;
721 }
722
723 return countSoFar;
724}
725
726template <class T, bool allowShortOperation>
727void Buffer<T, allowShortOperation>::addSegment(const T* buffer, size_t count) {
728 // Called with lock taken.
729 Segment* pNewSegment = new Segment();
730 pedigree_std::copy(pNewSegment->data, buffer, count);
731 pNewSegment->size = count;
732 m_Segments.pushBack(pNewSegment);
733}
734
735#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
736template <class T, bool allowShortOperation>
738 m_Lock.acquire();
739}
740
741template <class T, bool allowShortOperation>
743 m_Lock.release();
744}
745
746template <class T, bool allowShortOperation>
748 LockGuard<Mutex> guard(m_Lock);
749 return m_ActiveOperations;
750}
751#endif
752
753template class Buffer<uint8_t, false>;
754template class Buffer<uint8_t, true>;
755template class Buffer<char, false>;
756template class Buffer<char, true>;
757template class Buffer<size_t, false>;
758template class Buffer<size_t, true>;
size_t write(const T *buffer, size_t count, bool block=true)
Definition Buffer.cc:103
void notifyMonitorsLocked()
Definition Buffer.cc:672
size_t writeAtomic(const T *buffer, size_t count, bool block=true)
Definition Buffer.cc:227
size_t getSize()
Definition Buffer.cc:481
void monitor(Thread *pThread, Event *pEvent)
Definition Buffer.cc:594
void wipe()
Definition Buffer.cc:569
void disableReads()
Definition Buffer.cc:422
size_t getDataSize()
Definition Buffer.cc:470
void disableWrites()
Definition Buffer.cc:406
void cullMonitorTargets(Thread *pThread)
Definition Buffer.cc:620
void notifyMonitors()
Definition Buffer.cc:659
bool canRead(bool block)
Definition Buffer.cc:525
size_t writeLocked(const T *buffer, size_t count)
Definition Buffer.cc:679
bool canWrite(bool block)
Definition Buffer.cc:491
uint64_t readableGeneration() const
Definition Buffer.cc:559
size_t writeAvailable(const T *buffer, size_t count, bool atomic=false)
Definition Buffer.cc:201
void addSegment(const T *buffer, size_t count)
Definition Buffer.cc:727
bool enableWrites()
Definition Buffer.cc:438
bool enableReads()
Definition Buffer.cc:454
bool tryWrite(const T *buffer, size_t count)
Definition Buffer.cc:275
size_t read(T *buffer, size_t count, bool block=true)
Definition Buffer.cc:296
static bool mutexAcquired(Error error)
Definition Event.h:49
void disown()
Definition LockGuard.h:69
static ProcessorInformation & information()
size_t size
Definition Buffer.h:249
T data[m_SegmentSize]
Definition Buffer.h:243
size_t reader
Definition Buffer.h:246