|
27 | 27 | #include "Framework/InputSpec.h" |
28 | 28 | #include "Framework/Logger.h" |
29 | 29 | #include "Framework/OutputSpec.h" |
| 30 | +#include "Framework/RawDeviceService.h" |
30 | 31 | #include "Framework/Variant.h" |
31 | 32 | #include "../../../Algorithm/include/Algorithm/HeaderStack.h" |
32 | 33 | #include "Framework/OutputObjHeader.h" |
|
44 | 45 | #include <ROOT/RDataFrame.hxx> |
45 | 46 | #include <ROOT/RArrowDS.hxx> |
46 | 47 | #include <ROOT/RVec.hxx> |
| 48 | + |
| 49 | +#include <FairMQDevice.h> |
47 | 50 | #include <chrono> |
48 | 51 | #include <fstream> |
49 | 52 | #include <functional> |
@@ -506,6 +509,15 @@ DataProcessorSpec CommonDataProcessors::getDummySink(std::vector<InputSpec> cons |
506 | 509 | .algorithm = AlgorithmSpec{adaptStateful([](CallbackService& callbacks) { |
507 | 510 | auto dataConsumed = [](ServiceRegistry& services) { |
508 | 511 | services.get<DataProcessingStats>().consumedTimeframes++; |
| 512 | + auto device = services.get<RawDeviceService>().device(); |
| 513 | + auto channel = device->fChannels.find("metric-feedback"); |
| 514 | + if (channel != device->fChannels.end()) { |
| 515 | + FairMQMessagePtr payload(device->NewMessage()); |
| 516 | + int64_t* consumed = (int64_t*)malloc(sizeof(int64_t)); |
| 517 | + *consumed = services.get<DataProcessingStats>().consumedTimeframes; |
| 518 | + payload->Rebuild(consumed, sizeof(int64_t), nullptr, nullptr); |
| 519 | + channel->second[0].Send(payload); |
| 520 | + } |
509 | 521 | }; |
510 | 522 | callbacks.set(CallbackService::Id::DataConsumed, dataConsumed); |
511 | 523 |
|
|
0 commit comments