Add rx_streamdummyheader command (#1442)

* add rx_restream_stop command. This allows to send a ZMQ dummy header any time user wants to do so. For example this allows to pre-configure the ZMQ processing software before the acquisition begins. Therefore, the dummy header was adapted in order to contain the fields stored in the receiver.

* renamed command, changed inherit, moved commands to zmq related section

* update filename in restreamstop

* renamed helper function, sorted header fields alphabetically

* fixed fnametostream set

* renamed functions, add SetFileName method, format JSON parameters order

* added python bindings and formatting (does nothing really)

* renamed restream stop functions to stream dummy

* checkout .github files from ed8c885

* release notes

---------

Co-authored-by: Dhanya Thattil <dhanya.thattil@psi.ch>
This commit is contained in:
2026-08-03 09:27:53 +02:00
committed by mazzol_a
co-authored by maliakal_d
parent 504e5c0d58
commit faedb0acf1
25 changed files with 185 additions and 91 deletions
+4 -4
View File
@@ -167,7 +167,7 @@ int ClientInterface::functionTable(){
flist[F_GET_RECEIVER_STREAMING_PORT] = &ClientInterface::get_streaming_port;
flist[F_SET_RECEIVER_SILENT_MODE] = &ClientInterface::set_silent_mode;
flist[F_GET_RECEIVER_SILENT_MODE] = &ClientInterface::get_silent_mode;
flist[F_RESTREAM_STOP_FROM_RECEIVER] = &ClientInterface::restream_stop;
flist[F_STREAM_RX_DUMMY_HEADER_FROM_RECEIVER] = &ClientInterface::stream_rx_dummy_header;
flist[F_SET_ADDITIONAL_JSON_HEADER] = &ClientInterface::set_additional_json_header;
flist[F_GET_ADDITIONAL_JSON_HEADER] = &ClientInterface::get_additional_json_header;
flist[F_RECEIVER_UDP_SOCK_BUF_SIZE] = &ClientInterface::set_udp_socket_buffer_size;
@@ -1081,14 +1081,14 @@ int ClientInterface::get_silent_mode(Interface &socket) {
return socket.sendResult(retval);
}
int ClientInterface::restream_stop(Interface &socket) {
int ClientInterface::stream_rx_dummy_header(Interface &socket) {
verifyIdle(socket);
if (!impl()->getDataStreamEnable()) {
throw RuntimeError(
"Could not restream stop packet as data Streaming is disabled");
"Could not stream rx dummy header as data Streaming is disabled");
} else {
LOG(logDEBUG1) << "Restreaming stop";
impl()->restreamStop();
impl()->streamRxDummyHeader();
}
return socket.Send(OK);
}
+1 -1
View File
@@ -111,7 +111,7 @@ class ClientInterface : private virtual slsDetectorDefs {
int get_streaming_port(ServerInterface &socket);
int set_silent_mode(ServerInterface &socket);
int get_silent_mode(ServerInterface &socket);
int restream_stop(ServerInterface &socket);
int stream_rx_dummy_header(ServerInterface &socket);
int set_additional_json_header(ServerInterface &socket);
int get_additional_json_header(ServerInterface &socket);
int set_udp_socket_buffer_size(ServerInterface &socket);
+51 -39
View File
@@ -8,7 +8,6 @@
#include "DataStreamer.h"
#include "Fifo.h"
#include "GeneralData.h"
#include "sls/ZmqSocket.h"
#include "sls/sls_detector_exceptions.h"
#include <cerrno>
@@ -30,6 +29,10 @@ void DataStreamer::SetGeneralData(GeneralData *g) { generalData = g; }
void DataStreamer::SetFileIndex(uint64_t value) { fileIndex = value; }
void DataStreamer::SetFileName(const std::string &fname) {
fileNametoStream = fname;
}
void DataStreamer::SetNumberofPorts(xy np) { numPorts = np; }
void DataStreamer::SetFlipRows(bool fd) {
@@ -62,11 +65,10 @@ void DataStreamer::SetPortROI(ROI roi) {
}
}
void DataStreamer::ResetParametersforNewAcquisition(const std::string &fname) {
void DataStreamer::ResetParametersforNewAcquisition() {
StopRunning();
startedFlag = false;
firstIndex = 0;
fileNametoStream = fname;
}
void DataStreamer::RecordFirstIndex(uint64_t fnum, size_t firstImageIndex) {
@@ -152,8 +154,7 @@ void DataStreamer::ProcessAnImage(sls_detector_header header, size_t size,
uint64_t fnum = header.frameNumber;
LOG(logDEBUG1) << "DataStreamer " << index << ": fnum:" << fnum;
if (!SendDataHeader(header, size, generalData->nPixelsX,
generalData->nPixelsY)) {
if (!SendDataHeader(header, size)) {
LOG(logERROR) << "Could not send zmq header for fnum " << fnum
<< " and streamer " << index;
}
@@ -164,34 +165,47 @@ void DataStreamer::ProcessAnImage(sls_detector_header header, size_t size,
}
}
int DataStreamer::SendDummyHeader() {
zmqHeader DataStreamer::prepareRxZmqHeader() {
zmqHeader zHeader;
zHeader.data = false;
zHeader.jsonversion = SLS_DETECTOR_JSON_HEADER_VERSION;
// parameters coming from the receiver
zHeader.dynamicRange = generalData->dynamicRange;
zHeader.fileIndex = fileIndex;
zHeader.flipRows = static_cast<int>(flipRows);
zHeader.fname = fileNametoStream;
zHeader.imageSize = generalData->imageSize;
if (generalData->detType == GOTTHARD2 && index != 0) {
zHeader.imageSize = generalData->vetoImageSize;
}
zHeader.ndetx = numPorts.x;
zHeader.ndety = numPorts.y;
zHeader.npixelsx = generalData->nPixelsX;
zHeader.npixelsy = generalData->nPixelsY;
zHeader.quad = quadEnable;
// update local copy only if it was updated (to prevent locking each time)
if (isAdditionalJsonUpdated) {
std::lock_guard<std::mutex> lock(additionalJsonMutex);
localAdditionalJsonHeader = additionalJsonHeader;
isAdditionalJsonUpdated = false;
}
zHeader.addJsonHeader = localAdditionalJsonHeader;
zHeader.rx_roi = portRoi.getIntArray();
return zHeader;
}
int DataStreamer::SendDummyHeader() {
zmqHeader zHeader = prepareRxZmqHeader();
zHeader.data = false;
return zmqSocket->SendHeader(index, zHeader);
}
int DataStreamer::SendDataHeader(sls_detector_header header, uint32_t size,
uint32_t nx, uint32_t ny) {
zmqHeader zHeader;
int DataStreamer::SendDataHeader(sls_detector_header header, uint32_t size) {
zmqHeader zHeader = prepareRxZmqHeader();
zHeader.data = true;
zHeader.jsonversion = SLS_DETECTOR_JSON_HEADER_VERSION;
uint64_t frameIndex = header.frameNumber - firstIndex;
uint64_t acquisitionIndex = header.frameNumber;
zHeader.dynamicRange = generalData->dynamicRange;
zHeader.fileIndex = fileIndex;
zHeader.ndetx = numPorts.x;
zHeader.ndety = numPorts.y;
zHeader.npixelsx = nx;
zHeader.npixelsy = ny;
zHeader.imageSize = size;
zHeader.acqIndex = acquisitionIndex;
zHeader.frameIndex = frameIndex;
zHeader.progress =
100 * ((double)(frameIndex + 1) / (double)(nTotalFrames));
zHeader.fname = fileNametoStream;
// parameter from detector header
zHeader.frameNumber = header.frameNumber;
zHeader.expLength = header.expLength;
zHeader.packetNumber = header.packetNumber;
@@ -205,26 +219,24 @@ int DataStreamer::SendDataHeader(sls_detector_header header, uint32_t size,
zHeader.detSpec4 = header.detSpec4;
zHeader.detType = header.detType;
zHeader.version = header.version;
zHeader.flipRows = static_cast<int>(flipRows);
zHeader.quad = quadEnable;
// parameter derived from header and receiver
uint64_t acquisitionIndex = header.frameNumber;
uint64_t frameIndex = header.frameNumber - firstIndex;
zHeader.acqIndex = acquisitionIndex;
zHeader.completeImage =
(header.packetNumber < generalData->packetsPerFrame ? false : true);
// update local copy only if it was updated (to prevent locking each time)
if (isAdditionalJsonUpdated) {
std::lock_guard<std::mutex> lock(additionalJsonMutex);
localAdditionalJsonHeader = additionalJsonHeader;
isAdditionalJsonUpdated = false;
}
zHeader.addJsonHeader = localAdditionalJsonHeader;
zHeader.rx_roi = portRoi.getIntArray();
zHeader.frameIndex = frameIndex;
zHeader.imageSize = size;
zHeader.progress =
100 * ((double)(frameIndex + 1) / (double)(nTotalFrames));
return zmqSocket->SendHeader(index, zHeader);
}
void DataStreamer::RestreamStop() {
void DataStreamer::StreamRxDummyHeader() {
if (!SendDummyHeader()) {
throw RuntimeError("Could not restream Dummy Header via ZMQ for port " +
throw RuntimeError("Could not stream Dummy Header via ZMQ for port " +
std::to_string(zmqSocket->GetPortNumber()));
}
}
+6 -6
View File
@@ -10,6 +10,7 @@
*/
#include "ThreadObject.h"
#include "sls/ZmqSocket.h"
#include "sls/network_utils.h"
#include <map>
@@ -32,6 +33,7 @@ class DataStreamer : private virtual slsDetectorDefs, public ThreadObject {
void SetGeneralData(GeneralData *g);
void SetFileIndex(uint64_t value);
void SetFileName(const std::string &fname);
void SetNumberofPorts(xy np);
void SetFlipRows(bool fd);
void SetQuadEnable(bool value);
@@ -40,7 +42,7 @@ class DataStreamer : private virtual slsDetectorDefs, public ThreadObject {
SetAdditionalJsonHeader(const std::map<std::string, std::string> &json);
void SetPortROI(ROI roi);
void ResetParametersforNewAcquisition(const std::string &fname);
void ResetParametersforNewAcquisition();
/**
* Creates Zmq Sockets
* (throws an exception if it couldnt create zmq sockets)
@@ -49,7 +51,7 @@ class DataStreamer : private virtual slsDetectorDefs, public ThreadObject {
*/
void CreateZmqSockets(uint16_t port, int hwm);
void CloseZmqSocket();
void RestreamStop();
void StreamRxDummyHeader();
private:
/**
@@ -71,18 +73,16 @@ class DataStreamer : private virtual slsDetectorDefs, public ThreadObject {
*/
void ProcessAnImage(sls_detector_header header, size_t size, char *data);
zmqHeader prepareRxZmqHeader();
int SendDummyHeader();
/**
* Create and send Json Header
* @param rheader header of image
* @param size data size (could have been modified in call back)
* @param nx number of pixels in x dim
* @param ny number of pixels in y dim
* @returns 0 if error, else 1
*/
int SendDataHeader(sls_detector_header header, uint32_t size = 0,
uint32_t nx = 0, uint32_t ny = 0);
int SendDataHeader(sls_detector_header header, uint32_t size = 0);
static const std::string TypeName;
const GeneralData *generalData{nullptr};
+11 -6
View File
@@ -842,10 +842,13 @@ void Implementation::shutDownUDPSockets() {
it->ShutDownUDPSocket();
}
void Implementation::restreamStop() {
for (const auto &it : dataStreamer)
it->RestreamStop();
LOG(logINFO) << "Restreaming Dummy Header via ZMQ successful";
void Implementation::streamRxDummyHeader() {
std::string fnametostream = (filePath / fileName).string();
for (const auto &it : dataStreamer) {
it->SetFileName(fnametostream);
it->StreamRxDummyHeader();
}
LOG(logINFO) << "Streaming Dummy Header via ZMQ successful";
}
void Implementation::ResetParametersforNewAcquisition() {
@@ -856,8 +859,10 @@ void Implementation::ResetParametersforNewAcquisition() {
if (dataStreamEnable) {
std::string fnametostream = (filePath / fileName).string();
for (const auto &it : dataStreamer)
it->ResetParametersforNewAcquisition(fnametostream);
for (const auto &it : dataStreamer) {
it->ResetParametersforNewAcquisition();
it->SetFileName(fnametostream);
}
}
}
+1 -1
View File
@@ -104,7 +104,7 @@ class Implementation : private virtual slsDetectorDefs {
void stopReceiver();
void startReadout();
void shutDownUDPSockets();
void restreamStop();
void streamRxDummyHeader();
/**************************************************
* *