Project
Loading...
Searching...
No Matches
GPUWorkflowPipeline.cxx
Go to the documentation of this file.
1// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
2// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
3// All rights not expressly granted are reserved.
4//
5// This software is distributed under the terms of the GNU General Public
6// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
7//
8// In applying this license CERN does not waive the privileges and immunities
9// granted to it by virtue of its status as an Intergovernmental Organization
10// or submit itself to any jurisdiction.
11
14
17#include "GPUO2Interface.h"
18#include "GPUDataTypesIO.h"
19#include "GPUSettings.h"
20#include "GPUWorkflowInternal.h"
21
22#include "Framework/WorkflowSpec.h" // o2::framework::mergeInputs
30#include "Framework/Logger.h"
34
35#include <fairmq/Device.h>
36#include <fairmq/Channel.h>
37#include <fairmq/States.h>
38
39using namespace o2::framework;
40using namespace o2::header;
41using namespace o2::gpu;
42using namespace o2::base;
43using namespace o2::dataformats;
45
46namespace o2::gpu
47{
48
49static const std::string GPURecoWorkflowSpec_FMQCallbackKey = "GPURecoWorkflowSpec_FMQCallbackKey";
50
60
61void GPURecoWorkflowSpec::initPipeline(o2::framework::InitContext& ic)
62{
63 if (mSpecConfig.enableDoublePipeline == 1) {
64 mPipeline->fmqDevice = ic.services().get<RawDeviceService>().device();
65 mPipeline->fmqDevice->SubscribeToStateChange(GPURecoWorkflowSpec_FMQCallbackKey, [this](fair::mq::State s) { receiveFMQStateCallback(s); });
66 mPolicyOrder = [this](o2::framework::DataProcessingHeader::StartTime timeslice) {
67 std::unique_lock lk(mPipeline->completionPolicyMutex);
68 mPipeline->completionPolicyNotify.wait(lk, [pipeline = mPipeline.get()] { return pipeline->pipelineSenderTerminating || !pipeline->completionPolicyQueue.empty(); });
69 return !mPipeline->completionPolicyQueue.empty() && mPipeline->completionPolicyQueue.front() == timeslice;
70 };
71 mPipeline->receiveThread = std::thread([this]() { RunReceiveThread(); });
72 for (uint32_t i = 0; i < mPipeline->workers.size(); i++) {
73 mPipeline->workers[i].thread = std::thread([this, i]() { RunWorkerThread(i); });
74 }
75 }
76}
77
78void GPURecoWorkflowSpec::RunWorkerThread(int32_t id)
79{
80 LOG(debug) << "Running pipeline worker " << id;
81 auto& workerContext = mPipeline->workers[id];
82 while (!mPipeline->shouldTerminate) {
84 {
85 std::unique_lock lk(workerContext.inputQueueMutex);
86 workerContext.inputQueueNotify.wait(lk, [this, &workerContext]() { return mPipeline->shouldTerminate || !workerContext.inputQueue.empty(); });
87 if (workerContext.inputQueue.empty()) {
88 break;
89 }
90 context = workerContext.inputQueue.front();
91 workerContext.inputQueue.pop();
92 }
93 context->jobThreadIndex = id;
94 context->jobReturnValue = runMain(nullptr, context->jobPtrs, context->jobOutputRegions, id, context->jobInputUpdateCallback.get());
95 {
96 std::lock_guard lk(context->jobFinishedMutex);
97 context->jobFinished = true;
98 }
99 context->jobFinishedNotify.notify_one();
100 }
101}
102
103void GPURecoWorkflowSpec::enqueuePipelinedJob(GPUTrackingInOutPointers* ptrs, GPUInterfaceOutputs* outputRegions, GPURecoWorkflow_QueueObject* context, bool inputFinal)
104{
105 {
106 std::unique_lock lk(mPipeline->mayInjectMutex);
107 mPipeline->mayInjectCondition.wait(lk, [this, context]() { return mPipeline->mayInject && mPipeline->mayInjectTFId == context->mTFId; });
108 mPipeline->mayInjectTFId = mPipeline->mayInjectTFId + 1;
109 mPipeline->mayInject = false;
110 }
111 context->jobSubmitted = true;
112 context->jobInputFinal = inputFinal;
113 context->jobPtrs = ptrs;
114 context->jobOutputRegions = outputRegions;
115
116 context->jobInputUpdateCallback = std::make_unique<GPUInterfaceInputUpdate>();
117
118 if (!inputFinal) {
119 context->jobInputUpdateCallback->callback = [context, this](GPUTrackingInOutPointers*& data, GPUInterfaceOutputs*& outputs) -> int32_t {
120 std::unique_lock lk(context->jobInputFinalMutex);
121 context->jobInputFinalNotify.wait(lk, [context, this]() { return context->jobInputFinal || mPipeline->pipelineAbort; });
122 if (mPipeline->pipelineAbort) {
123 return 1;
124 }
125 data = context->jobPtrs;
126 outputs = context->jobOutputRegions;
127 return 0;
128 };
129 }
130 context->jobInputUpdateCallback->notifyCallback = [this]() {
131 {
132 std::lock_guard lk(mPipeline->mayInjectMutex);
133 mPipeline->mayInject = true;
134 }
135 mPipeline->mayInjectCondition.notify_one();
136 };
137
138 mNextThreadIndex = (mNextThreadIndex + 1) % 2;
139
140 {
141 std::lock_guard lk(mPipeline->workers[mNextThreadIndex].inputQueueMutex);
142 mPipeline->workers[mNextThreadIndex].inputQueue.emplace(context);
143 }
144 mPipeline->workers[mNextThreadIndex].inputQueueNotify.notify_one();
145}
146
147void GPURecoWorkflowSpec::finalizeInputPipelinedJob(GPUTrackingInOutPointers* ptrs, GPUInterfaceOutputs* outputRegions, GPURecoWorkflow_QueueObject* context)
148{
149 {
150 std::lock_guard lk(context->jobInputFinalMutex);
151 context->jobPtrs = ptrs;
152 context->jobOutputRegions = outputRegions;
153 context->jobInputFinal = true;
154 }
155 context->jobInputFinalNotify.notify_one();
156}
157
158int32_t GPURecoWorkflowSpec::handlePipeline(ProcessingContext& pc, GPUTrackingInOutPointers& ptrs, GPURecoWorkflowSpec_TPCZSBuffers& tpcZSmeta, o2::gpu::GPUTrackingInOutZS& tpcZS, std::unique_ptr<GPURecoWorkflow_QueueObject>& context)
159{
160 mPipeline->runStarted = true;
161 mPipeline->stateNotify.notify_all();
162
163 auto* device = pc.services().get<RawDeviceService>().device();
164 const auto& tinfo = pc.services().get<o2::framework::TimingInfo>();
165 if (mSpecConfig.enableDoublePipeline == 1) {
166 std::unique_lock lk(mPipeline->queueMutex);
167 mPipeline->queueNotify.wait(lk, [this] { return !mPipeline->pipelineQueue.empty(); });
168 context = std::move(mPipeline->pipelineQueue.front());
169 mPipeline->pipelineQueue.pop();
170 lk.unlock();
171
172 if (context->timeSliceId != tinfo.timeslice) {
173 LOG(fatal) << "Prepare message for incorrect time frame received, time frames seem out of sync";
174 }
175
176 tpcZSmeta = std::move(context->tpcZSmeta);
177 tpcZS = context->tpcZS;
178 ptrs.tpcZS = &tpcZS;
179
180 {
181 std::lock_guard lk(mPipeline->completionPolicyMutex);
182 if (mPipeline->completionPolicyQueue.empty() || mPipeline->completionPolicyQueue.front() != tinfo.timeslice) {
183 LOG(fatal) << "Time frame processed does not equal the timeframe at the top of the queue, time frames seem out of sync";
184 }
185 mPipeline->completionPolicyQueue.pop();
186 }
187 } else if (mSpecConfig.enableDoublePipeline == 2) {
188 auto prepareDummyMessage = pc.outputs().make<DataAllocator::UninitializedVector<char>>(Output{gDataOriginGPU, "PIPELINEPREPARE", 0}, 0u);
189
190 size_t ptrsTotal = 0;
191 const void* firstPtr = nullptr;
192 for (uint32_t i = 0; i < GPUTrackingInOutZS::NSECTORS; i++) {
193 for (uint32_t j = 0; j < GPUTrackingInOutZS::NENDPOINTS; j++) {
194 if (firstPtr == nullptr && ptrs.tpcZS->sector[i].count[j]) {
195 firstPtr = ptrs.tpcZS->sector[i].zsPtr[j][0];
196 }
197 ptrsTotal += ptrs.tpcZS->sector[i].count[j];
198 }
199 }
200
201 size_t prepareBufferSize = sizeof(pipelinePrepareMessage) + ptrsTotal * sizeof(size_t) * 4;
202 fair::mq::MessagePtr payload(device->NewMessage());
203 payload->Rebuild(prepareBufferSize, fair::mq::Alignment(sizeof(size_t)));
204 auto* messageBuffer = (size_t*)payload->GetData();
205 pipelinePrepareMessage& preMessage = *(pipelinePrepareMessage*)messageBuffer;
206 preMessage.magicWord = preMessage.MAGIC_WORD;
207 preMessage.timeSliceId = tinfo.timeslice;
208 preMessage.pointersTotal = ptrsTotal;
209 preMessage.flagEndOfStream = false;
210 memcpy((void*)&preMessage.tfSettings, (const void*)ptrs.settingsTF, sizeof(preMessage.tfSettings));
211
212 size_t* ptrBuffer = messageBuffer + sizeof(preMessage) / sizeof(size_t);
213 size_t ptrsCopied = 0;
214 int32_t lastRegion = -1;
215 for (uint32_t i = 0; i < GPUTrackingInOutZS::NSECTORS; i++) {
216 for (uint32_t j = 0; j < GPUTrackingInOutZS::NENDPOINTS; j++) {
217 preMessage.pointerCounts[i][j] = ptrs.tpcZS->sector[i].count[j];
218 for (uint32_t k = 0; k < ptrs.tpcZS->sector[i].count[j]; k++) {
219 const void* curPtr = ptrs.tpcZS->sector[i].zsPtr[j][k];
220 bool regionFound = lastRegion != -1 && (size_t)curPtr >= (size_t)mRegionInfos[lastRegion].ptr && (size_t)curPtr < (size_t)mRegionInfos[lastRegion].ptr + mRegionInfos[lastRegion].size;
221 if (!regionFound) {
222 for (uint32_t l = 0; l < mRegionInfos.size(); l++) {
223 if ((size_t)curPtr >= (size_t)mRegionInfos[l].ptr && (size_t)curPtr < (size_t)mRegionInfos[l].ptr + mRegionInfos[l].size) {
224 lastRegion = l;
225 regionFound = true;
226 break;
227 }
228 }
229 }
230 if (!regionFound) {
231 LOG(fatal) << "Found a TPC ZS pointer outside of shared memory";
232 }
233 ptrBuffer[ptrsCopied + k] = (size_t)curPtr - (size_t)mRegionInfos[lastRegion].ptr;
234 ptrBuffer[ptrsTotal + ptrsCopied + k] = ptrs.tpcZS->sector[i].nZSPtr[j][k];
235 ptrBuffer[2 * ptrsTotal + ptrsCopied + k] = mRegionInfos[lastRegion].managed;
236 ptrBuffer[3 * ptrsTotal + ptrsCopied + k] = mRegionInfos[lastRegion].id;
237 }
238 ptrsCopied += ptrs.tpcZS->sector[i].count[j];
239 }
240 }
241
242 auto channel = device->GetChannels().find("gpu-prepare-channel");
243 LOG(info) << "Sending gpu-reco-workflow prepare message of size " << prepareBufferSize;
244 channel->second[0].Send(payload);
245 return 2;
246 }
247 return 0;
248}
249
250void GPURecoWorkflowSpec::handlePipelineEndOfStream(EndOfStreamContext& ec)
251{
252 if (mSpecConfig.enableDoublePipeline == 1) {
253 mPipeline->endOfStreamDplReceived = true;
254 mPipeline->stateNotify.notify_all();
255 }
256 if (mSpecConfig.enableDoublePipeline == 2) {
257 auto* device = ec.services().get<RawDeviceService>().device();
258 fair::mq::MessagePtr payload(device->NewMessage());
259 payload->Rebuild(sizeof(pipelinePrepareMessage), fair::mq::Alignment(alignof(pipelinePrepareMessage)));
260 auto* preMessage = (pipelinePrepareMessage*)payload->GetData();
261 new (preMessage) pipelinePrepareMessage;
262 preMessage->flagEndOfStream = true;
263 auto channel = device->GetChannels().find("gpu-prepare-channel");
264 LOG(info) << "Sending end-of-stream message over out-of-bands channel";
265 channel->second[0].Send(payload);
266 }
267}
268
269void GPURecoWorkflowSpec::handlePipelineStop()
270{
271 if (mSpecConfig.enableDoublePipeline == 1) {
272 {
273 std::unique_lock lk(mPipeline->queueMutex);
274 mPipeline->pipelineAbort = mPipeline->pipelineQueue.size();
275 }
276 if (mPipeline->pipelineAbort) {
277 mPipeline->pipelineQueue.front()->jobInputFinalNotify.notify_one();
278 mGPUReco->DrainPipeline();
279 {
280 std::unique_lock lk(mPipeline->queueMutex);
281 mPipeline->pipelineQueue = {};
282 }
283 {
284 std::lock_guard lk(mPipeline->completionPolicyMutex);
285 mPipeline->completionPolicyQueue = {};
286 }
287 mPipeline->pipelineAbort = false;
288 {
289 std::lock_guard lk(mPipeline->stateMutex);
290 mPipeline->endOfStreamAsyncWaiting = false;
291 mPipeline->mNTFReceived = 0;
292 mPipeline->runStarted = false;
293 }
294 }
295 {
296 std::unique_lock lk(mPipeline->mayInjectMutex);
297 mPipeline->mayInjectTFId = 0;
298 }
299 }
300}
301
302void GPURecoWorkflowSpec::receiveFMQStateCallback(fair::mq::State newState)
303{
304 {
305 std::lock_guard lk(mPipeline->stateMutex);
306 if (mPipeline->fmqState != fair::mq::State::Running && newState == fair::mq::State::Running) {
307 mPipeline->endOfStreamAsyncWaiting = true;
308 mPipeline->endOfStreamDplReceived = false;
309 }
310 mPipeline->fmqPreviousState = mPipeline->fmqState;
311 mPipeline->fmqState = newState;
312 }
313 mPipeline->stateNotify.notify_all();
314 {
315 std::lock_guard lk(mPipeline->receiveMutex);
316 if (newState == fair::mq::State::Exiting) {
317 mPipeline->fmqDevice->UnsubscribeFromStateChange(GPURecoWorkflowSpec_FMQCallbackKey);
318 }
319 }
320}
321
322void GPURecoWorkflowSpec::RunReceiveThread()
323{
324 auto* device = mPipeline->fmqDevice;
325 while (!mPipeline->shouldTerminate) {
326 bool received = false;
327 int32_t recvTimeot = 1000;
328 fair::mq::MessagePtr msg;
329 LOG(debug) << "Waiting for out of band message";
330 auto shouldReceive = [this]() { return ((mPipeline->fmqState == fair::mq::State::Running || (mPipeline->fmqState == fair::mq::State::Ready && mPipeline->fmqPreviousState == fair::mq::State::Running)) && mPipeline->endOfStreamAsyncWaiting); };
331 do {
332 {
333 std::unique_lock lk(mPipeline->stateMutex);
334 mPipeline->stateNotify.wait(lk, [this, shouldReceive]() { return shouldReceive() || mPipeline->shouldTerminate; }); // Do not check mPipeline->fmqDevice->NewStatePending() since we wait for EndOfStream!
335 }
336 if (mPipeline->shouldTerminate) {
337 break;
338 }
339 try {
340 do {
341 std::unique_lock lk(mPipeline->receiveMutex);
342 if (!shouldReceive()) {
343 break;
344 }
345 msg = device->NewMessageFor("gpu-prepare-channel", 0, 0);
346 received = device->Receive(msg, "gpu-prepare-channel", 0, recvTimeot) > 0;
347 } while (!received && !mPipeline->shouldTerminate);
348 } catch (...) {
349 usleep(1000000);
350 }
351 } while (!received && !mPipeline->shouldTerminate);
352 if (mPipeline->shouldTerminate) {
353 break;
354 }
355 if (msg->GetSize() < sizeof(pipelinePrepareMessage)) {
356 LOG(fatal) << "Received prepare message of invalid size " << msg->GetSize() << " < " << sizeof(pipelinePrepareMessage);
357 }
358 const pipelinePrepareMessage* m = (const pipelinePrepareMessage*)msg->GetData();
359 if (m->magicWord != m->MAGIC_WORD) {
360 LOG(fatal) << "Prepare message corrupted, invalid magic word";
361 }
362 if (m->flagEndOfStream) {
363 LOG(info) << "Received end-of-stream from out-of-band channel";
364 {
365 std::lock_guard lk(mPipeline->stateMutex);
366 mPipeline->endOfStreamAsyncWaiting = false;
367 mPipeline->mNTFReceived = 0;
368 mPipeline->runStarted = false;
369 }
370 mPipeline->stateNotify.notify_all();
371 continue;
372 }
373
374 {
375 std::lock_guard lk(mPipeline->completionPolicyMutex);
376 mPipeline->completionPolicyQueue.emplace(m->timeSliceId);
377 }
378 mPipeline->completionPolicyNotify.notify_one();
379
380 {
381 std::unique_lock lk(mPipeline->stateMutex);
382 mPipeline->stateNotify.wait(lk, [this]() { return (mPipeline->runStarted && mPipeline->endOfStreamAsyncWaiting) || mPipeline->shouldTerminate; });
383 if (!mPipeline->runStarted) {
384 continue;
385 }
386 }
387
388 auto context = std::make_unique<GPURecoWorkflow_QueueObject>();
389 context->timeSliceId = m->timeSliceId;
390 context->tfSettings = m->tfSettings;
391
392 size_t ptrsCopied = 0;
393 size_t* ptrBuffer = (size_t*)msg->GetData() + sizeof(pipelinePrepareMessage) / sizeof(size_t);
394 context->tpcZSmeta.Pointers[0][0].resize(m->pointersTotal);
395 context->tpcZSmeta.Sizes[0][0].resize(m->pointersTotal);
396 int32_t lastRegion = -1;
397 for (uint32_t i = 0; i < GPUTrackingInOutZS::NSECTORS; i++) {
398 for (uint32_t j = 0; j < GPUTrackingInOutZS::NENDPOINTS; j++) {
399 context->tpcZS.sector[i].count[j] = m->pointerCounts[i][j];
400 for (uint32_t k = 0; k < context->tpcZS.sector[i].count[j]; k++) {
401 bool regionManaged = ptrBuffer[2 * m->pointersTotal + ptrsCopied + k];
402 size_t regionId = ptrBuffer[3 * m->pointersTotal + ptrsCopied + k];
403 bool regionFound = lastRegion != -1 && mRegionInfos[lastRegion].managed == regionManaged && mRegionInfos[lastRegion].id == regionId;
404 if (!regionFound) {
405 for (uint32_t l = 0; l < mRegionInfos.size(); l++) {
406 if (mRegionInfos[l].managed == regionManaged && mRegionInfos[l].id == regionId) {
407 lastRegion = l;
408 regionFound = true;
409 break;
410 }
411 }
412 }
413 if (!regionFound) {
414 LOG(fatal) << "Received ZS Ptr for SHM region (managed " << (int32_t)regionManaged << ", id " << regionId << "), which was not registered for us";
415 }
416 context->tpcZSmeta.Pointers[0][0][ptrsCopied + k] = (void*)(ptrBuffer[ptrsCopied + k] + (size_t)mRegionInfos[lastRegion].ptr);
417 context->tpcZSmeta.Sizes[0][0][ptrsCopied + k] = ptrBuffer[m->pointersTotal + ptrsCopied + k];
418 }
419 context->tpcZS.sector[i].zsPtr[j] = context->tpcZSmeta.Pointers[0][0].data() + ptrsCopied;
420 context->tpcZS.sector[i].nZSPtr[j] = context->tpcZSmeta.Sizes[0][0].data() + ptrsCopied;
421 ptrsCopied += context->tpcZS.sector[i].count[j];
422 }
423 }
424 context->ptrs.tpcZS = &context->tpcZS;
425 context->ptrs.settingsTF = &context->tfSettings;
426 context->mTFId = mPipeline->mNTFReceived;
427 if (mPipeline->mNTFReceived++ >= mPipeline->workers.size()) { // Do not inject the first workers.size() TFs, since we need a first round of calib updates from DPL before starting
428 enqueuePipelinedJob(&context->ptrs, nullptr, context.get(), false);
429 }
430 {
431 std::lock_guard lk(mPipeline->queueMutex);
432 mPipeline->pipelineQueue.emplace(std::move(context));
433 }
434 mPipeline->queueNotify.notify_one();
435 }
436 mPipeline->pipelineSenderTerminating = true;
437 mPipeline->completionPolicyNotify.notify_one();
438}
439
440void GPURecoWorkflowSpec::ExitPipeline()
441{
442 if (mSpecConfig.enableDoublePipeline == 1 && mPipeline->fmqDevice) {
443 mPipeline->fmqDevice = nullptr;
444 mPipeline->shouldTerminate = true;
445 mPipeline->stateNotify.notify_all();
446 for (uint32_t i = 0; i < mPipeline->workers.size(); i++) {
447 mPipeline->workers[i].inputQueueNotify.notify_one();
448 }
449 if (mPipeline->receiveThread.joinable()) {
450 mPipeline->receiveThread.join();
451 }
452 for (uint32_t i = 0; i < mPipeline->workers.size(); i++) {
453 if (mPipeline->workers[i].thread.joinable()) {
454 mPipeline->workers[i].thread.join();
455 }
456 }
457 }
458}
459
460} // namespace o2::gpu
std::ostringstream debug
int32_t i
A helper class to iteratate over all parts of all input routes.
uint32_t j
Definition RawData.h:0
Type wrappers for enfording a specific serialization method.
TBranch * ptr
decltype(auto) make(const Output &spec, Args... args)
ServiceRegistryRef services()
Definition InitContext.h:34
DataAllocator & outputs()
The data allocator is used to allocate memory for the output data.
ServiceRegistryRef services()
The services registry associated with this processing context.
o2::framework::Outputs outputs()
const GLfloat * m
Definition glcorearb.h:4066
GLsizeiptr size
Definition glcorearb.h:659
GLboolean * data
Definition glcorearb.h:298
GLuint id
Definition glcorearb.h:650
constexpr o2::header::DataOrigin gDataOriginGPU
Definition DataHeader.h:593
Definition of a container to keep/associate and arbitrary number of labels associated to an index wit...
Defining ITS Vertex explicitly as messageable.
Definition Cartesian.h:288
O2 data header classes and API, v0.1.
Definition DetID.h:49
const GPUTrackingInOutZS * tpcZS
GPUTrackingInOutZSSector sector[NSECTORS]
static constexpr uint32_t NSECTORS
static constexpr uint32_t NENDPOINTS
std::unique_ptr< GPUInterfaceInputUpdate > jobInputUpdateCallback
size_t pointerCounts[GPUTrackingInOutZS::NSECTORS][GPUTrackingInOutZS::NENDPOINTS]
DataProcessingHeader::StartTime timeSliceId
LOG(info)<< "Compressed in "<< sw.CpuTime()<< " s"
uint64_t const void const *restrict const msg
Definition x9.h:153