The Pedigree Project 0.1
RequestQueue.h
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#ifndef REQUEST_QUEUE_H
21#define REQUEST_QUEUE_H
22#include "pedigree/kernel/Atomic.h"
23#include "pedigree/kernel/compiler.h"
24#include "pedigree/kernel/process/Mutex.h"
25#include "pedigree/kernel/process/SchedulerWorkerWake.h"
26#include "pedigree/kernel/process/WaitQueue.h"
27#include "pedigree/kernel/processor/state_forward.h"
28#include "pedigree/kernel/processor/types.h"
29#include "pedigree/kernel/utilities/IntrusiveMpscQueue.h"
30#include "pedigree/kernel/utilities/StaticString.h"
31#include "pedigree/kernel/utilities/String.h"
32
33#include <config.h>
34#if THREADS
35#include "pedigree/kernel/machine/TimerHandler.h"
36#endif
37
38class Thread;
39struct ThreadPlacement;
40class Timer;
42
43#define REQUEST_QUEUE_NUM_PRIORITIES 4
44
52class EXPORTED_PUBLIC RequestQueue {
53 public:
55
56 enum class OverrunStatus {
57 Clear,
58 Armed,
59 Stalled,
60 Overloaded,
61 };
62
63#if THREADS && HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
65 OverrunStatus sampleOverrunForTest();
66
67 using HostedSmokeHook = void (*)(void*);
68
69 void setAfterPreallocatedAdmissionHookForTest(HostedSmokeHook hook, void* context) {
70 m_AfterPreallocatedAdmissionHook = hook;
71 m_AfterPreallocatedAdmissionContext = context;
72 }
73
74 void setAfterIntakeExchangeHookForTest(HostedSmokeHook hook, void* context) {
75 m_AfterIntakeExchangeHook = hook;
76 m_AfterIntakeExchangeContext = context;
77 }
78
79 size_t workerTransientRetriesForTest() {
80 return m_WorkerTransientRetries.value();
81 }
82
83 size_t guardedTransientRetriesForTest() {
84 return m_GuardedTransientRetries.value();
85 }
86
87 size_t publisherDrainRetriesForTest() {
88 return m_PublisherDrainRetries.value();
89 }
90#endif
91
92#if THREADS && PEDIGREE_CONCURRENCY_SMOKE_TESTS
93 using ConcurrencySmokeHook = void (*)(void*);
94
95 void setAfterPreallocatedClaimHookForTest(ConcurrencySmokeHook hook, void* context) {
96 m_AfterPreallocatedClaimHook = hook;
97 m_AfterPreallocatedClaimContext = context;
98 }
99
100#endif
101
102#if THREADS
103 private:
105 friend class RequestQueue;
106
108 : m_LastQueueSize(0),
109 m_LastProgressGeneration(0),
110 m_Tick(0),
111 m_HasBacklogBaseline(false),
112 queue(0) {}
113
114 private:
115 virtual void timer(uint64_t delta);
116 OverrunStatus sample(size_t& lastSize, size_t& currentSize);
117 void resetBaselineLocked();
118
119 size_t m_LastQueueSize;
120 size_t m_LastProgressGeneration;
121 uint64_t m_Tick;
122 bool m_HasBacklogBaseline;
123
124 RequestQueue* queue;
125 };
126#endif
127
128 protected:
129 class Request;
130
131 struct IntakeNode {
132 explicit IntakeNode(Request* request = nullptr) : next(nullptr), owner(request) {}
133
134 IntakeNode* next;
135 Request* owner;
136 };
137
139 class Request {
140 public:
141 Request(size_t requestPriority, bool asynchronous, uint64_t requestP1, uint64_t requestP2,
142 uint64_t requestP3, uint64_t requestP4, uint64_t requestP5, uint64_t requestP6,
143 uint64_t requestP7, uint64_t requestP8,
144 PreallocatedRequest* preallocatedOwner = nullptr)
145 : p1(requestP1),
146 p2(requestP2),
147 p3(requestP3),
148 p4(requestP4),
149 p5(requestP5),
150 p6(requestP6),
151 p7(requestP7),
152 p8(requestP8),
153 m_ReturnValue(0),
154#if THREADS
155 m_Completion(),
156 m_References(asynchronous ? 1 : 2),
157#endif
158 m_Next(nullptr),
159 m_Intake(this),
160 m_Priority(requestPriority),
161 m_Asynchronous(asynchronous),
162 m_Rejected(false),
163 m_Completed(false),
164 m_pPreallocatedOwner(preallocatedOwner) {
165 }
166
167 uint64_t p1, p2, p3, p4, p5, p6, p7, p8;
168
169 private:
170 friend class PreallocatedRequest;
171 friend class RequestQueue;
172
173 ~Request() = default;
174
175 uint64_t m_ReturnValue;
176#if THREADS
177 WaitQueue m_Completion;
178 Atomic<size_t> m_References;
179#endif
180 Request* m_Next;
181 IntakeNode m_Intake;
182 size_t m_Priority;
183 bool m_Asynchronous;
184 bool m_Rejected;
185 bool m_Completed;
186 PreallocatedRequest* m_pPreallocatedOwner;
187
188 Request(const Request&);
189 void operator=(const Request&);
190 };
191
193 public:
194 IntakeLane() : m_Stub(), m_Queue(m_Stub) {}
195
196 IntakeNode m_Stub;
198
199 private:
200 IntakeLane(const IntakeLane&) = delete;
201 IntakeLane& operator=(const IntakeLane&) = delete;
202 };
203
204 public:
224 class EXPORTED_PUBLIC PreallocatedRequest {
225 public:
226 using ReleaseCallback = void (*)(void*);
227
229 PreallocatedRequest(ReleaseCallback releaseCallback, void* releaseContext);
231
232 bool isAvailable() const;
233
234 private:
235 friend class RequestQueue;
236 NOT_COPYABLE_OR_ASSIGNABLE(PreallocatedRequest);
237
238 enum State {
239 Idle,
240 Claimed,
241 Published,
242 Releasing,
243 };
244
245 Request m_Request;
246 Atomic<size_t> m_State;
247 Atomic<size_t> m_ReleaseDepth;
248 ReleaseCallback m_ReleaseCallback;
249 void* m_ReleaseContext;
250 };
251
252 enum class PreallocatedPublishResult {
253 Accepted,
254 // Another accepted publication, or a claim that cannot roll back,
255 // already owns this token.
256 TokenBusy,
257 QueueStopped,
258 // Retained for source and module ABI compatibility. Preallocated
259 // publications do not consume allocation admission.
260 QueueFull,
261 InvalidPriority,
262 };
263
265 RequestQueue(const String& name);
266 virtual ~RequestQueue();
267
268 // Action to perform when a duplicate request is found in the queue.
269 enum ActionOnDuplicate {
270 // Block waiting for it to complete, and return its return value.
271 Block,
272 // Ignore the duplicate and create a new request.
273 NewRequest,
274 // Return immediately, ignoring the result from the request.
275 ReturnImmediately
276 };
277
278 enum class LifecycleState {
279 Stopped,
280 Accepting,
281 Stopping,
282 Destroyed,
283 };
284
286 virtual void initialise();
287
296 virtual void destroy();
297
308 MUST_USE_RESULT uint64_t addRequest(size_t priority, uint64_t p1 = 0, uint64_t p2 = 0,
309 uint64_t p3 = 0, uint64_t p4 = 0, uint64_t p5 = 0,
310 uint64_t p6 = 0, uint64_t p7 = 0, uint64_t p8 = 0);
311
314 MUST_USE_RESULT uint64_t addRequest(size_t priority, ActionOnDuplicate action, uint64_t p1 = 0,
315 uint64_t p2 = 0, uint64_t p3 = 0, uint64_t p4 = 0,
316 uint64_t p5 = 0, uint64_t p6 = 0, uint64_t p7 = 0,
317 uint64_t p8 = 0);
318
328 uint64_t addAsyncRequest(size_t priority, uint64_t p1 = 0, uint64_t p2 = 0, uint64_t p3 = 0,
329 uint64_t p4 = 0, uint64_t p5 = 0, uint64_t p6 = 0, uint64_t p7 = 0,
330 uint64_t p8 = 0);
331
339 MUST_USE_RESULT bool publishAsyncRequest(size_t priority, uint64_t p1 = 0, uint64_t p2 = 0,
340 uint64_t p3 = 0, uint64_t p4 = 0, uint64_t p5 = 0,
341 uint64_t p6 = 0, uint64_t p7 = 0, uint64_t p8 = 0);
342
344 bool canWaitForCompletion();
345
348 MUST_USE_RESULT bool waitForPreallocated(PreallocatedRequest& request);
349
366 MUST_USE_RESULT PreallocatedPublishResult publishPreallocated(PreallocatedRequest& request,
367 size_t priority, uint64_t p1 = 0,
368 uint64_t p2 = 0, uint64_t p3 = 0,
369 uint64_t p4 = 0, uint64_t p5 = 0,
370 uint64_t p6 = 0, uint64_t p7 = 0,
371 uint64_t p8 = 0);
372
380 MUST_USE_RESULT PreallocatedPublishResult republishPreallocatedWhileReleasing(
381 PreallocatedRequest& request, size_t priority, uint64_t p1 = 0, uint64_t p2 = 0,
382 uint64_t p3 = 0, uint64_t p4 = 0, uint64_t p5 = 0, uint64_t p6 = 0, uint64_t p7 = 0,
383 uint64_t p8 = 0);
384
391 MUST_USE_RESULT bool halt();
392
396 MUST_USE_RESULT bool resume();
397
399 LifecycleState getLifecycleState();
400
406 bool drain();
407
408 protected:
409 virtual bool workerPlacement(ThreadPlacement&) const {
410 return false;
411 }
412
416 virtual uint64_t executeRequest(uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4, uint64_t p5,
417 uint64_t p6, uint64_t p7, uint64_t p8) = 0;
418
420 void operator=(const RequestQueue&);
421
429 virtual bool compareRequests(const Request& a, const Request& b) {
430 return false;
431 }
432
444 virtual void cancelRequest(const Request& request) {}
445
447 static int trampoline(void* p);
448
450 int work();
451
452 enum class NextRequestResult {
453 Item,
454 Empty,
455 Retry,
456 };
457
459 NextRequestResult getNextRequest(Request*& request);
460
462 bool drainIntakeLocked(size_t priority);
463
465 void publishRequest(Request* request);
466
468 Request* findDuplicate(const Request& request);
469
471 static void completeRequest(Request* request, uint64_t returnValue, bool rejected);
472
474 void discardRequest(Request* request);
475
476 uint64_t addAsyncRequestInternal(size_t priority, uint64_t p1, uint64_t p2, uint64_t p3,
477 uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8);
478
479 PreallocatedPublishResult publishPreallocatedRequest(PreallocatedRequest& request,
480 PreallocatedRequest::State availableState,
481 size_t priority, uint64_t p1, uint64_t p2,
482 uint64_t p3, uint64_t p4, uint64_t p5,
483 uint64_t p6, uint64_t p7, uint64_t p8);
484
486 void releasePreallocatedRequest(Request* request);
487
489 void invokeCancelRequest(const Request& request);
490
492 bool callbackActiveOnCurrentThread() const;
493
494#if THREADS
496 static void retainRequest(Request* request);
497 void releaseRequest(Request* request);
498 uint64_t waitForRequest(Request* request);
499
501 bool startWorker();
502 bool stopWorker();
503
504 void closePreallocatedAdmission();
505 void waitForPreallocatedPublishers();
506#endif
507
508 static constexpr size_t PublicationClosed = static_cast<size_t>(1) << ((sizeof(size_t) * 8) - 1);
509 static constexpr size_t PublicationCountMask = ~PublicationClosed;
510
512 IntakeLane m_IntakeLanes[REQUEST_QUEUE_NUM_PRIORITIES];
513
515 Request* m_pRequestQueue[REQUEST_QUEUE_NUM_PRIORITIES];
516 Request* m_pRequestQueueTail[REQUEST_QUEUE_NUM_PRIORITIES];
517
520
523
524#if THREADS
527
530 WaitQueue m_WorkerWaiters;
531
532 Thread* m_pThread;
533
536 SchedulerWorkerWake m_WorkerWake;
537
540
543
546
547 RequestQueueOverrunChecker m_OverrunChecker;
548 Timer* m_pOverrunTimer;
549
552
553#if HOSTED && PEDIGREE_HOSTED_SMOKE_TESTS
554 HostedSmokeHook m_AfterPreallocatedAdmissionHook;
555 void* m_AfterPreallocatedAdmissionContext;
556 HostedSmokeHook m_AfterIntakeExchangeHook;
557 void* m_AfterIntakeExchangeContext;
558 Atomic<size_t> m_WorkerTransientRetries;
559 Atomic<size_t> m_GuardedTransientRetries;
560 Atomic<size_t> m_PublisherDrainRetries;
561#endif
562#if PEDIGREE_CONCURRENCY_SMOKE_TESTS
563 ConcurrencySmokeHook m_AfterPreallocatedClaimHook;
564 void* m_AfterPreallocatedClaimContext;
565#endif
566#endif
567
570 Atomic<size_t> m_nAsyncRequests;
571
574
577
578 NormalStaticString m_Name;
579};
580
581#endif
Definition Mutex.h:56
Atomic< PerProcessorScheduler * > m_pWorkerScheduler
virtual bool compareRequests(const Request &a, const Request &b)
size_t m_WorkerProgressGeneration
WaitQueue m_RequestQueueWaiters
Atomic< size_t > m_bWorkerActive
Atomic< size_t > m_State
Mutex m_LifecycleMutex
Atomic< size_t > m_bWorkerReady
size_t m_nMaxAsyncRequests
Atomic< size_t > m_nTotalRequests
virtual void cancelRequest(const Request &request)
virtual uint64_t executeRequest(uint64_t p1, uint64_t p2, uint64_t p3, uint64_t p4, uint64_t p5, uint64_t p6, uint64_t p7, uint64_t p8)=0
Atomic< size_t > m_nActiveRequests
Request * m_pActiveRequest
Atomic< size_t > m_PublicationState