Skip to content

Commit 68019b1

Browse files
committed
GPU Workflow: Fix API usage of fairmq message
1 parent 2bd5701 commit 68019b1

1 file changed

Lines changed: 10 additions & 9 deletions

File tree

GPU/Workflow/src/GPUWorkflowPipeline.cxx

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -199,15 +199,17 @@ int32_t GPURecoWorkflowSpec::handlePipeline(ProcessingContext& pc, GPUTrackingIn
199199
}
200200

201201
size_t prepareBufferSize = sizeof(pipelinePrepareMessage) + ptrsTotal * sizeof(size_t) * 4;
202-
std::vector<size_t> messageBuffer(prepareBufferSize / sizeof(size_t));
203-
pipelinePrepareMessage& preMessage = *(pipelinePrepareMessage*)messageBuffer.data();
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;
204206
preMessage.magicWord = preMessage.MAGIC_WORD;
205207
preMessage.timeSliceId = tinfo.timeslice;
206208
preMessage.pointersTotal = ptrsTotal;
207209
preMessage.flagEndOfStream = false;
208210
memcpy((void*)&preMessage.tfSettings, (const void*)ptrs.settingsTF, sizeof(preMessage.tfSettings));
209211

210-
size_t* ptrBuffer = messageBuffer.data() + sizeof(preMessage) / sizeof(size_t);
212+
size_t* ptrBuffer = messageBuffer + sizeof(preMessage) / sizeof(size_t);
211213
size_t ptrsCopied = 0;
212214
int32_t lastRegion = -1;
213215
for (uint32_t i = 0; i < GPUTrackingInOutZS::NSECTORS; i++) {
@@ -238,9 +240,7 @@ int32_t GPURecoWorkflowSpec::handlePipeline(ProcessingContext& pc, GPUTrackingIn
238240
}
239241

240242
auto channel = device->GetChannels().find("gpu-prepare-channel");
241-
fair::mq::MessagePtr payload(device->NewMessage());
242243
LOG(info) << "Sending gpu-reco-workflow prepare message of size " << prepareBufferSize;
243-
payload->Rebuild(messageBuffer.data(), prepareBufferSize, nullptr, nullptr);
244244
channel->second[0].Send(payload);
245245
return 2;
246246
}
@@ -255,12 +255,13 @@ void GPURecoWorkflowSpec::handlePipelineEndOfStream(EndOfStreamContext& ec)
255255
}
256256
if (mSpecConfig.enableDoublePipeline == 2) {
257257
auto* device = ec.services().get<RawDeviceService>().device();
258-
pipelinePrepareMessage preMessage;
259-
preMessage.flagEndOfStream = true;
260-
auto channel = device->GetChannels().find("gpu-prepare-channel");
261258
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");
262264
LOG(info) << "Sending end-of-stream message over out-of-bands channel";
263-
payload->Rebuild(&preMessage, sizeof(preMessage), nullptr, nullptr);
264265
channel->second[0].Send(payload);
265266
}
266267
}

0 commit comments

Comments
 (0)