Skip to content

Commit 3400eb4

Browse files
committed
GPU: Delay only processing of async message not receiving, to fill completionPolicyQueue in time and do not deadlock
1 parent 55d3067 commit 3400eb4

2 files changed

Lines changed: 17 additions & 11 deletions

File tree

GPU/Workflow/src/GPUWorkflowInternal.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -84,7 +84,7 @@ struct GPURecoWorkflowSpec_PipelineInternals {
8484
std::condition_variable queueNotify;
8585

8686
std::queue<o2::framework::DataProcessingHeader::StartTime> completionPolicyQueue;
87-
bool pipelineSenderTerminating = false;
87+
volatile bool pipelineSenderTerminating = false;
8888
std::mutex completionPolicyMutex;
8989
std::condition_variable completionPolicyNotify;
9090

GPU/Workflow/src/GPUWorkflowPipeline.cxx

Lines changed: 16 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -283,7 +283,7 @@ void GPURecoWorkflowSpec::RunReceiveThread()
283283
do {
284284
{
285285
std::unique_lock lk(mPipeline->stateMutex);
286-
mPipeline->stateNotify.wait(lk, [this]() { return (mPipeline->runStarted && mPipeline->fmqState == fair::mq::State::Running && !mPipeline->endOfStreamReceived) || mPipeline->shouldTerminate; }); // Do not check mPipeline->fmqDevice->NewStatePending() since we wait for EndOfStream!
286+
mPipeline->stateNotify.wait(lk, [this]() { return (mPipeline->fmqState == fair::mq::State::Running && !mPipeline->endOfStreamReceived) || mPipeline->shouldTerminate; }); // Do not check mPipeline->fmqDevice->NewStatePending() since we wait for EndOfStream!
287287
}
288288
if (mPipeline->shouldTerminate) {
289289
break;
@@ -316,6 +316,20 @@ void GPURecoWorkflowSpec::RunReceiveThread()
316316
continue;
317317
}
318318

319+
{
320+
std::lock_guard lk(mPipeline->completionPolicyMutex);
321+
mPipeline->completionPolicyQueue.emplace(m->timeSliceId);
322+
}
323+
mPipeline->completionPolicyNotify.notify_one();
324+
325+
{
326+
std::unique_lock lk(mPipeline->stateMutex);
327+
mPipeline->stateNotify.wait(lk, [this]() { return (mPipeline->runStarted && !mPipeline->endOfStreamReceived) || mPipeline->shouldTerminate; });
328+
if (!mPipeline->runStarted) {
329+
continue;
330+
}
331+
}
332+
319333
auto context = std::make_unique<GPURecoWorkflow_QueueObject>();
320334
context->timeSliceId = m->timeSliceId;
321335
context->tfSettings = m->tfSettings;
@@ -358,21 +372,13 @@ void GPURecoWorkflowSpec::RunReceiveThread()
358372
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
359373
enqueuePipelinedJob(&context->ptrs, nullptr, context.get(), false);
360374
}
361-
{
362-
std::lock_guard lk(mPipeline->completionPolicyMutex);
363-
mPipeline->completionPolicyQueue.emplace(m->timeSliceId);
364-
}
365-
mPipeline->completionPolicyNotify.notify_one();
366375
{
367376
std::lock_guard lk(mPipeline->queueMutex);
368377
mPipeline->pipelineQueue.emplace(std::move(context));
369378
}
370379
mPipeline->queueNotify.notify_one();
371380
}
372-
{
373-
std::lock_guard lk(mPipeline->completionPolicyMutex);
374-
mPipeline->pipelineSenderTerminating = true;
375-
}
381+
mPipeline->pipelineSenderTerminating = true;
376382
mPipeline->completionPolicyNotify.notify_one();
377383
}
378384

0 commit comments

Comments
 (0)