Files
Jungfraujoch/tools/jfjoch_stream2_replay.cpp
leonarski_fandClaude Opus 5.5 8e134c405d tools: run the broker on a stored dataset, without a detector
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
2026-10-08 15:52:22 +02:00

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;
}