The Pedigree Project 0.1
ProducerConsumer.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/utilities/ProducerConsumer.h"
21
22#if PRODUCERCONSUMER_ASYNCHRONOUS
23#include "pedigree/kernel/LockGuard.h"
24#include "pedigree/kernel/Log.h"
25#include "pedigree/kernel/utilities/pocketknife.h"
26#endif
27
28ProducerConsumer::ProducerConsumer() = default;
29
30ProducerConsumer::~ProducerConsumer() {
31#if PRODUCERCONSUMER_ASYNCHRONOUS
32 if (m_pThreadHandle) {
33 FATAL(
34 "ProducerConsumer destroyed before its most-derived destructor "
35 "stopped the worker.");
36 }
37
38 // Tasks queued on an object that was never initialised are still owned by
39 // the base class and cannot have reached virtual consume().
40 for (auto it : m_Tasks) {
41 delete it;
42 }
43#endif
44}
45
47#if PRODUCERCONSUMER_ASYNCHRONOUS
48 m_Lock.acquire();
49 if (m_Destroyed) {
50 m_Lock.release();
51 return;
52 }
53
54 m_Destroyed = true;
55 m_Running = false;
56 m_Condition.signal();
57 void* threadHandle = m_pThreadHandle;
58 m_pThreadHandle = nullptr;
59 m_Lock.release();
60
61 if (threadHandle) {
62 if (!pocketknife::attachToForCompletion(threadHandle)) {
63 FATAL("ProducerConsumer could not join its worker during teardown.");
64 }
65 }
66
67 // Clean up tasks that didn't get executed.
68 for (auto it : m_Tasks) {
69 delete it;
70 }
71 m_Tasks.clear();
72#endif
73}
74
75bool ProducerConsumer::initialise() {
76#if PRODUCERCONSUMER_ASYNCHRONOUS
77 LockGuard<Mutex> guard(m_Lock);
78
79 if (m_Destroyed) {
80 return false;
81 }
82
83 if (m_Running) {
84 return true;
85 }
86
87 m_Running = true;
88 m_pThreadHandle = pocketknife::runConcurrentlyAttached(thread, this);
89 if (!m_pThreadHandle) {
90 m_Running = false;
91 return false;
92 }
93
94 return true;
95#else
96 return true;
97#endif
98}
99
100void ProducerConsumer::produce(uint64_t p0, uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4,
101 uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8) {
102#if PRODUCERCONSUMER_ASYNCHRONOUS
103 Task* task = new Task;
104 task->p0 = p0;
105 task->p1 = p1;
106 task->p2 = p2;
107 task->p3 = p3;
108 task->p4 = p4;
109 task->p5 = p5;
110 task->p6 = p6;
111 task->p7 = p7;
112 task->p8 = p8;
113
114 m_Lock.acquire();
115 if (m_Destroyed) {
116 m_Lock.release();
117 delete task;
118 return;
119 }
120 m_Tasks.pushBack(task);
121 m_Condition.signal();
122 m_Lock.release();
123#else
124 consume(p0, p1, p2, p3, p4, p5, p6, p7, p8);
125#endif
126}
127
128void ProducerConsumer::consumerThread() {
129 LockGuard<Mutex> guard(m_Lock);
130
131 while (true) {
132 while (m_Running && !m_Tasks.size()) {
133 ConditionVariable::Error error = ConditionVariable::NoError;
134 if (!m_Condition.wait(m_Lock, error)) {
136 guard.disown();
137 }
138 return;
139 }
140 }
141
142 if (!m_Running) {
143 break;
144 }
145
146 Task* task = m_Tasks.popFront();
147
148 // Don't hold lock while we actually perform the consume operation.
149 m_Lock.release();
150
151 consume(task->p0, task->p1, task->p2, task->p3, task->p4, task->p5, task->p6, task->p7,
152 task->p8);
153
154 delete task;
155
156 m_Lock.acquire();
157 }
158}
159
160int ProducerConsumer::thread(void* p) {
161 ProducerConsumer* pc = reinterpret_cast<ProducerConsumer*>(p);
162 pc->consumerThread();
163
164 return 0;
165}
MUST_USE_RESULT bool wait(Mutex &mutex, Time::Timestamp &timeout, Error &error, WaitQueue::StackDiscardCleanup onStackDiscard=nullptr, void *stackDiscardContext=nullptr)
static bool mutexAcquired(Error error)
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
EXPORTED_PUBLIC bool attachToForCompletion(void *handle)
EXPORTED_PUBLIC void * runConcurrentlyAttached(int(*func)(void *), void *param)