The Pedigree Project 0.1
Pipe.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 "Pipe.h"
21#include "pedigree/kernel/LockGuard.h"
22#include "pedigree/kernel/process/Mutex.h"
23#include "pedigree/kernel/process/Process.h"
24#include "pedigree/kernel/process/Thread.h"
25#include "pedigree/kernel/processor/Processor.h"
26#include "pedigree/kernel/processor/ProcessorInformation.h"
27#include "pedigree/kernel/utilities/ZombieQueue.h"
28#include "pedigree/kernel/utilities/new"
29
30class Filesystem;
31
32class ZombiePipe : public ZombieObject {
33 public:
34 ZombiePipe(Pipe* pPipe) : m_pPipe(pPipe) {}
35 virtual ~ZombiePipe();
36
37 private:
38 Pipe* m_pPipe;
39};
40
41ZombiePipe::~ZombiePipe() {
42 NOTICE("ZombiePipe: freeing " << m_pPipe);
43 delete m_pPipe;
44}
45
47 : File(),
48 m_bIsAnonymous(true),
49 m_bIsEOF(false),
50 m_Buffer(bufferChanged, this),
51 m_ReaderCondition(),
52 m_WriteGeneration(0),
53 m_ErrorGeneration(0),
54 m_HangupGeneration(0),
55 m_nLifetimePins(0),
56 m_bRetirementQueued(false) {
57 if constexpr (VERBOSE_KERNEL) {
58 NOTICE("Pipe: new anonymous pipe " << reinterpret_cast<uintptr_t>(this));
59 }
60}
61
62Pipe::Pipe(const String& name, Time::Timestamp accessedTime, Time::Timestamp modifiedTime,
63 Time::Timestamp creationTime, uintptr_t inode, Filesystem* pFs, size_t size,
64 File* pParent, bool bIsAnonymous)
65 : File(name, accessedTime, modifiedTime, creationTime, inode, pFs, size, pParent),
66 m_bIsAnonymous(bIsAnonymous),
67 m_bIsEOF(false),
68 m_Buffer(bufferChanged, this),
69 m_ReaderCondition(),
70 m_WriteGeneration(0),
71 m_ErrorGeneration(0),
72 m_HangupGeneration(0),
73 m_nLifetimePins(0),
74 m_bRetirementQueued(false) {
75 if constexpr (VERBOSE_KERNEL) {
76 NOTICE("Pipe: new " << (bIsAnonymous ? "anonymous" : "named") << " pipe " << Hex << this);
77 }
78}
79
81 // ensure anything else in the critical section can finish before we clean
82 // up fully
83 // this is useful for cases where ZombieQueue destroys us before we get a
84 // chance to actually return from decreaseRefCount (which accesses the lock)
85 m_Lock.acquire();
86 m_Lock.release();
87}
88
89void Pipe::bufferChanged(void* context) {
90 static_cast<Pipe*>(context)->dataChanged();
91}
92
94 return m_Buffer.getDataSize();
95}
96
97int Pipe::select(bool bWriting, int timeout) {
98 if (bWriting) {
99 return m_Buffer.canWrite(timeout > 0) ? 1 : 0;
100 } else {
101 return m_Buffer.canRead(timeout > 0) ? 1 : 0;
102 }
103}
104
105ReadyMask Pipe::queryReady(bool reading, bool writing) {
106 LockGuard<Mutex> guard(m_Lock);
107 ReadyMask ready = ReadyNone;
108
109 if (!m_nWriters) {
110 ready |= ReadyHangup;
111 }
112 if (!m_nReaders) {
113 ready |= ReadyError;
114 }
115
116 if (reading) {
117 if (m_Buffer.canRead(false)) {
118 ready |= ReadyRead;
119 }
120 }
121
122 if (writing) {
123 // A closed read end is immediately writable from poll/epoll's point of
124 // view even though the subsequent write fails with EPIPE/SIGPIPE.
125 if (!m_nReaders || m_Buffer.canWrite(false)) {
126 ready |= ReadyWrite;
127 }
128 }
129
130 return ready;
131}
132
134 LockGuard<Mutex> guard(m_Lock);
135 ReadinessGenerations generations;
136 generations.read = m_Buffer.readableGeneration();
137 generations.write = m_Buffer.writableGeneration() + m_WriteGeneration;
138 generations.error = m_ErrorGeneration;
139 generations.hangup = m_HangupGeneration;
140 return generations;
141}
142
143uint64_t Pipe::readBytewise(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock) {
144 // Need to read what's left in the pipe then EOF if there's no more readers!
145 {
146 LockGuard<Mutex> guard(m_Lock);
147 if (m_nWriters == 0) {
148 bCanBlock = false;
149 }
150 }
151
152 uint8_t* pBuf = reinterpret_cast<uint8_t*>(buffer);
153 return m_Buffer.read(pBuf, size, bCanBlock);
154}
155
156uint64_t Pipe::writeBytewise(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock) {
157 {
158 LockGuard<Mutex> guard(m_Lock);
159 if (m_nReaders == 0) {
160 // no more readers, abort the write
161 return 0;
162 }
163 }
164
165 uint8_t* pBuf = reinterpret_cast<uint8_t*>(buffer);
166 return size <= PIPE_BUF_MAX ? m_Buffer.writeAtomic(pBuf, size, bCanBlock)
167 : m_Buffer.write(pBuf, size, bCanBlock);
168}
169
170bool Pipe::isPipe() const {
171 return getName().length() == 0 || m_bIsAnonymous;
172}
173
174bool Pipe::isFifo() const {
175 return getName().length() > 0 && !m_bIsAnonymous;
176}
177
178void Pipe::increaseRefCount(bool bIsWriter) {
179 {
180 LockGuard<Mutex> guard(m_Lock);
181
182 if (bIsWriter) {
183 // A reader can still own unread bytes across the last writer's close.
184 m_Buffer.enableWrites();
185 m_nWriters++;
186 } else {
187 // A reader is now present so we can enable reads if they weren't.
188 m_Buffer.enableReads();
189 m_nReaders++;
190
191 // The predicate is "at least one reader", so one arrival satisfies
192 // every writer currently blocked in open().
194 }
195 }
196
197 dataChanged();
198}
199
200void Pipe::decreaseRefCount(bool bIsWriter) {
201 // Make sure only one thread decreases the refcount at a time. This is
202 // important as we add ourselves to the ZombieQueue if the refcount ticks
203 // to zero. Getting pre-empted by another thread that also decreases the
204 // refcount between the decrement and the check for zero may mean the pipe
205 // is added to the ZombieQueue twice, which causes a double free.
206 bool bDataChanged = false;
207 bool queueRetirement = false;
208 {
209 LockGuard<Mutex> guard(m_Lock);
210
211 if (m_nReaders == 0 && m_nWriters == 0) {
212 // Refcount is already zero - don't decrement! (also, bad.)
213 ERROR("Pipe: decreasing refcount when refcount is already zero.");
214 return;
215 }
216
217 if (bIsWriter) {
218 m_nWriters--;
219 if (m_nWriters == 0) {
220 ++m_HangupGeneration;
221 // Wakes up readers waiting as they won't be able to be woken
222 // by new bytes being written anymore.
223 m_Buffer.disableWrites();
224 bDataChanged = true;
225 }
226 } else {
227 const bool wasReadyForWrite = !m_nReaders || m_Buffer.canWrite(false);
228 m_nReaders--;
229 if (m_nReaders == 0) {
230 if (!wasReadyForWrite) {
231 ++m_WriteGeneration;
232 }
233 ++m_ErrorGeneration;
234 // Wake up any writers that were waiting for space - no more
235 // readers (EOF condition, pipe other end has left).
236 m_Buffer.disableReads();
237 bDataChanged = true;
238 }
239 }
240
241 if (!m_nReaders && !m_nWriters) {
242 // Named FIFO storage survives its final open description, unlike its data.
243 m_Buffer.wipe();
244 }
245 queueRetirement = shouldQueueRetirementLocked();
246 if (queueRetirement) {
247 bDataChanged = false;
248 }
249 }
250
251 if (queueRetirement) {
252 size_t pid = Processor::information().getCurrentThread()->getParent()->getId();
253 if constexpr (VERBOSE_KERNEL) {
254 NOTICE("Adding pipe [" << pid << "] " << this << " to ZombieQueue");
255 }
256 ZombieQueue::instance().addObject(new ZombiePipe(this));
257 return;
258 }
259
260 if (bDataChanged) {
261 dataChanged();
262 }
263}
264
266 if (!m_bIsAnonymous) {
268 }
269
270 LockGuard<Mutex> guard(m_Lock);
272 return false;
273 }
275 return true;
276}
277
279 if (!m_bIsAnonymous) {
281 return;
282 }
283
284 bool queueRetirement = false;
285 {
286 LockGuard<Mutex> guard(m_Lock);
287 assert(m_nLifetimePins);
289 queueRetirement = shouldQueueRetirementLocked();
290 }
291
292 if (queueRetirement) {
293 size_t pid = Processor::information().getCurrentThread()->getParent()->getId();
294 if constexpr (VERBOSE_KERNEL) {
295 NOTICE("Adding pipe [" << pid << "] " << this << " to ZombieQueue");
296 }
297 ZombieQueue::instance().addObject(new ZombiePipe(this));
298 }
299}
300
301bool Pipe::shouldQueueRetirementLocked() {
302 if (!m_bIsAnonymous || m_bRetirementQueued || m_nLifetimePins || m_nReaders || m_nWriters) {
303 return false;
304 }
305
306 m_bRetirementQueued = true;
307 return true;
308}
309
311 LockGuard<Mutex> guard(m_Lock);
312 return m_nReaders;
313}
314
316 LockGuard<Mutex> guard(m_Lock);
317 return m_nWriters;
318}
319
320bool Pipe::waitForReader(bool bCanBlock) {
321 m_Lock.acquire();
322 while (!m_nReaders) {
323 if (!bCanBlock) {
324 m_Lock.release();
325 return false;
326 }
327
328 ConditionVariable::Error error = ConditionVariable::NoError;
329 if (!m_ReaderCondition.wait(m_Lock, error)) {
331 m_Lock.release();
332 }
333 return false;
334 }
335 }
336 m_Lock.release();
337 return true;
338}
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 File.h:75
virtual bool retainVfsReference()
Definition File.cc:910
String getName() const
Definition File.cc:782
virtual void releaseVfsReference()
Definition File.cc:914
void dataChanged()
Definition File.cc:1452
Definition Pipe.h:36
ReadinessGenerations readinessGenerations() override
Definition Pipe.cc:133
virtual bool isFifo() const
Definition Pipe.cc:174
void releaseVfsReference() override
Definition Pipe.cc:278
size_t getWriterCount()
Definition Pipe.cc:315
bool m_bIsAnonymous
Definition Pipe.h:163
bool waitForReader(bool bCanBlock)
Definition Pipe.cc:320
bool retainVfsReference() override
Definition Pipe.cc:265
Pipe()
Definition Pipe.cc:46
size_t m_nLifetimePins
Definition Pipe.h:179
size_t readableBytes()
Definition Pipe.cc:93
bool m_bRetirementQueued
Definition Pipe.h:182
virtual bool isPipe() const
Definition Pipe.cc:170
virtual uint64_t readBytewise(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock=true)
Definition Pipe.cc:143
virtual void decreaseRefCount(bool bIsWriter)
Definition Pipe.cc:200
ReadyMask queryReady(bool reading, bool writing) override
Definition Pipe.cc:105
virtual int select(bool bWriting=false, int timeout=0)
Definition Pipe.cc:97
size_t getReaderCount()
Definition Pipe.cc:310
PipeBuffer m_Buffer
Definition Pipe.h:169
virtual uint64_t writeBytewise(uint64_t location, uint64_t size, uintptr_t buffer, bool bCanBlock=true)
Definition Pipe.cc:156
virtual ~Pipe()
Definition Pipe.cc:80
ConditionVariable m_ReaderCondition
Definition Pipe.h:172
static ProcessorInformation & information()
void release(size_t n=1)
Definition Semaphore.cc:549
bool acquire(size_t n=1, size_t timeoutSecs=0, size_t timeoutUsecs=0)
Definition Semaphore.cc:355
@ Hex
Definition Log.h:124