For demonstrating and testing a collection followed while it is written: - jfjoch_stream2_replay replays a stored dataset as a DECTRIS detector's stream2 (ZeroMQ PUSH on the stream2 port): start message, every image at a given rate, end message. - fake_simplon.py stands in for the detector's SIMPLON control interface: answers the configuration and status requests the broker makes, accepts every command, reads back idle. It needs port 80. Not installed. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
97 lines
3.8 KiB
C++
97 lines
3.8 KiB
C++
// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
|
// SPDX-License-Identifier: GPL-3.0-only
|
|
|
|
// Replays a stored dataset as the image stream of a DECTRIS detector: binds a ZeroMQ PUSH socket where
|
|
// the broker's DECTRIS receiver connects (tcp://*:31001) and sends the start message, every image at
|
|
// the given rate, and the end message. With a stand-in for the detector's control interface, the broker
|
|
// then runs a real collection - real geometry, real diffraction - with no detector. For demonstrations
|
|
// and tests only.
|
|
|
|
#include <chrono>
|
|
#include <iostream>
|
|
#include <thread>
|
|
#include <getopt.h>
|
|
|
|
#include "../common/Logger.h"
|
|
#include "../common/ZMQWrappers.h"
|
|
#include "../frame_serialize/CBORStream2Serializer.h"
|
|
#include "../reader/JFJochHDF5Reader.h"
|
|
|
|
void print_usage() {
|
|
std::cout << "Usage: jfjoch_stream2_replay [-r <images per s>] [-n <images>] [-a <address>] <input master>" << std::endl;
|
|
std::cout << " -r images sent per second (default 100)" << std::endl;
|
|
std::cout << " -n images to send (default all)" << std::endl;
|
|
std::cout << " -a ZeroMQ address to bind (default tcp://0.0.0.0:31001, the DECTRIS stream2 port)" << std::endl;
|
|
}
|
|
|
|
int main(int argc, char **argv) {
|
|
Logger logger("jfjoch_stream2_replay");
|
|
RegisterHDF5Filter();
|
|
|
|
double rate = 100.0;
|
|
std::optional<int64_t> images;
|
|
std::string address = "tcp://0.0.0.0:31001";
|
|
|
|
int opt;
|
|
while ((opt = getopt(argc, argv, "r:n:a:")) != -1) {
|
|
switch (opt) {
|
|
case 'r': rate = atof(optarg); break;
|
|
case 'n': images = atoll(optarg); break;
|
|
case 'a': address = optarg; break;
|
|
default: print_usage(); return EXIT_FAILURE;
|
|
}
|
|
}
|
|
if (optind + 1 != argc || rate <= 0.0) {
|
|
print_usage();
|
|
return EXIT_FAILURE;
|
|
}
|
|
|
|
JFJochHDF5Reader reader;
|
|
reader.ReadFile(argv[optind]);
|
|
const auto dataset = reader.GetDataset();
|
|
const int64_t n = std::min<int64_t>(images.value_or(reader.GetNumberOfImages()), reader.GetNumberOfImages());
|
|
|
|
// What a DECTRIS detector states about itself in its start message; the broker takes the
|
|
// geometry of the collection from its own /start.
|
|
DiffractionExperiment x(dataset->experiment);
|
|
x.ImagesPerTrigger(n).NumTriggers(1);
|
|
StartMessage start;
|
|
x.FillMessage(start);
|
|
const auto stored = reader.GetStoredPixelFormat();
|
|
start.bit_depth_image = stored.bit_depth;
|
|
start.pixel_signed = stored.is_signed;
|
|
start.number_of_images = n;
|
|
|
|
ZMQSocket socket(ZMQSocketType::Push);
|
|
socket.SendWaterMark(1000);
|
|
socket.Bind(address);
|
|
logger.Info("Bound {}; sending {} images at {} Hz", address, n, rate);
|
|
|
|
std::vector<uint8_t> buffer(256 * 1024 * 1024);
|
|
CBORStream2Serializer serializer(buffer.data(), buffer.size());
|
|
serializer.SerializeSequenceStart(start);
|
|
socket.Send(buffer.data(), serializer.GetBufferSize());
|
|
|
|
const auto t0 = std::chrono::steady_clock::now();
|
|
JFJochReaderRawImage raw;
|
|
for (int64_t i = 0; i < n; i++) {
|
|
std::this_thread::sleep_until(t0 + std::chrono::duration<double>(static_cast<double>(i) / rate));
|
|
reader.ReadRawImage(i, raw);
|
|
DataMessage message{};
|
|
message.image = raw.image;
|
|
message.number = i;
|
|
serializer.SerializeImage(message);
|
|
socket.Send(buffer.data(), serializer.GetBufferSize());
|
|
}
|
|
|
|
EndMessage end{};
|
|
end.max_image_number = n;
|
|
serializer.SerializeSequenceEnd(end);
|
|
socket.Send(buffer.data(), serializer.GetBufferSize());
|
|
logger.Info("Sent {} images in {:.1f} s", n,
|
|
std::chrono::duration<double>(std::chrono::steady_clock::now() - t0).count());
|
|
// Let the end message leave before the socket closes.
|
|
std::this_thread::sleep_for(std::chrono::seconds(2));
|
|
return EXIT_SUCCESS;
|
|
}
|