Files
Jungfraujoch/tests/LiveCollectionTest.cpp
leonarski_fandClaude Opus 5.5 2981aedd3e Broker: a collection can be followed from /start, its start message comes in the event stream
With a DECTRIS detector the broker starts the writers - and so learned the collection's start
message - only when the detector's own start message arrived, after /start had returned. A reader
started right after /start found no collection (404) and had to wait for, and race with, the
detector.

- /start registers the collection in LiveCollection (Expect) as it is accepted, so GET /live/events
  is accepted from then on. The image pusher fills in the start message when it starts the writers
  (Begin) - at once on the FPGA path, when the detector's start message arrives on the DECTRIS path.
- The event stream's first event is now `start`, carrying that start message (CBOR, base64) with
  run number, file prefix, image count and images per file; it replaces `collection`. Until it is
  known the stream only keeps alive. file and end follow as before.
- However the measurement thread ends, it ends a collection the pusher did not end (EndIfOpen), so a
  cancel or a failure before the detector streamed reaches the reader as `end` with an error. The
  wait is event-driven: it ends with the start message or with the measurement.
- GET /live/start.cbor is removed: the stream carries the start message, one route less.
- rugnux takes the start message from the stream (BrokerFeed::WaitForStart), so it can be started
  right after /start; files reported before it attached the reader are replayed to it.

Tests: LiveCollection_ExpectedAtStart; BrokerFeed_StartFromEventStream (an event stream whose start
message comes late, then a file and the end).

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
2026-10-08 17:03:09 +02:00

380 lines
15 KiB
C++

// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include <catch2/catch_all.hpp>
#include <httplib.h>
#include <base64/Base64.h>
#include <filesystem>
#include <future>
#include <thread>
#include "../common/DiffractionExperiment.h"
#include "../compression/JFJochCompressor.h"
#include "../image_puller/TCPImagePuller.h"
#include "../image_pusher/LiveCollection.h"
#include "../image_pusher/TCPStreamPusher.h"
#include "../reader/JFJochHDF5Reader.h"
#include "../reader/JFJochStartMessageReader.h"
#include "../rugnux/BrokerFeed.h"
#include "../writer/FileWriter.h"
#include "../writer/StreamWriter.h"
// Following a collection while it is written: the writer reports each closed data file, the broker
// tracks them (LiveCollection), and rugnux reads the images from the start message and the files.
namespace {
const std::string live_dir = "live_collection_test";
DiffractionExperiment LiveExperiment(const std::string &prefix) {
DiffractionExperiment x(DetJF(1));
x.FilePrefix(prefix).ImagesPerTrigger(7).OverwriteExistingFiles(true);
x.BitDepthImage(16).ImagesPerFile(2).SetFileWriterFormat(FileWriterFormat::NXmxVDS).PixelSigned(true);
x.Compression(CompressionAlgorithm::BSHUF_ZSTD);
x.BeamX_pxl(512.5).BeamY_pxl(256.25).DetectorDistance_mm(120.0).IncidentEnergy_keV(12.4);
x.Goniometer(GoniometerAxis("omega", 10.0f, 0.25f, Coord(-1, 0, 0), {}));
return x;
}
std::vector<uint8_t> Image(const DiffractionExperiment &x, int i) {
std::vector<int16_t> image(x.GetPixelsNum());
for (size_t p = 0; p < image.size(); p++)
image[p] = static_cast<int16_t>((p * 13 + 17 * i + 3) % 20000);
JFJochBitShuffleCompressor compressor(CompressionAlgorithm::BSHUF_ZSTD);
return compressor.Compress(image);
}
DataMessage Message(const DiffractionExperiment &x, const std::vector<uint8_t> &compressed, int i) {
DataMessage message{};
message.image = CompressedImage(compressed, x.GetXPixelsNum(), x.GetYPixelsNum(),
CompressedImageMode::Int16, CompressionAlgorithm::BSHUF_ZSTD);
message.number = i;
return message;
}
}
TEST_CASE("LiveCollection_Lifecycle", "[Live]") {
LiveCollection live;
CHECK(live.Get().generation == 0);
CHECK(!live.StartCBOR());
StartMessage start{};
start.file_prefix = "sub/run";
start.images_per_file = 100;
start.number_of_images = 250;
start.run_number = 7;
live.Begin(start, {1, 2, 3});
CHECK(live.StartCBOR() == std::vector<uint8_t>{1, 2, 3});
auto s = live.Wait(1, true, 0, std::chrono::milliseconds(0));
CHECK(s.generation == 1);
CHECK(s.run_number == 7);
CHECK(s.files.empty());
// A waiter wakes when a file is reported.
auto waiter = std::async(std::launch::async, [&] { return live.Wait(1, true, 0, std::chrono::seconds(10)); });
std::this_thread::sleep_for(std::chrono::milliseconds(50));
live.FileClosed(1, 100);
s = waiter.get();
REQUIRE(s.files.size() == 1);
CHECK(s.files[0].file == "sub/run_data_000002.h5");
CHECK(s.files[0].first_image == 100);
CHECK(s.files[0].image_count == 100);
live.FileClosed(2, 50);
live.End(250);
s = live.Get();
CHECK(s.files.size() == 2);
CHECK(s.ended);
CHECK(s.end_images == 250);
live.FileClosed(0, 100); // nothing after the end
CHECK(live.Get().files.size() == 2);
// One follower at a time.
CHECK(live.AcquireFollower());
CHECK_FALSE(live.AcquireFollower());
live.ReleaseFollower();
CHECK(live.AcquireFollower());
// A new collection starts from nothing.
live.Begin(start, {});
s = live.Get();
CHECK(s.generation == 2);
CHECK(s.files.empty());
CHECK_FALSE(s.ended);
}
TEST_CASE("LiveCollection_ExpectedAtStart", "[Live]") {
// /start announces the collection; the start message follows when the image pusher starts.
LiveCollection live;
live.Expect();
auto s = live.Get();
CHECK(s.generation == 1);
CHECK_FALSE(s.started);
CHECK(!live.StartCBOR());
// A follower waiting for the start message wakes when it comes, in the same collection.
auto waiter = std::async(std::launch::async, [&] { return live.Wait(1, false, 0, std::chrono::seconds(10)); });
std::this_thread::sleep_for(std::chrono::milliseconds(50));
StartMessage start{};
start.file_prefix = "run";
start.images_per_file = 10;
start.run_number = 3;
live.Begin(start, {4, 5});
s = waiter.get();
CHECK(s.generation == 1);
CHECK(s.started);
CHECK(s.run_number == 3);
CHECK(live.StartCBOR() == std::vector<uint8_t>{4, 5});
// The pusher's end stands; the state machine's afterwards changes nothing.
live.End(20);
live.EndIfOpen("measurement over");
CHECK(live.Get().end_images == 20);
CHECK(live.Get().end_error.empty());
// A measurement that ends before the start message arrives ends the collection with its error.
live.Expect();
waiter = std::async(std::launch::async, [&] { return live.Wait(2, false, 0, std::chrono::seconds(10)); });
std::this_thread::sleep_for(std::chrono::milliseconds(50));
live.EndIfOpen("measurement over");
s = waiter.get();
CHECK(s.generation == 2);
CHECK_FALSE(s.started);
CHECK(s.ended);
CHECK(s.end_error == "measurement over");
CHECK(!live.StartCBOR());
}
TEST_CASE("ServerSentEventParser", "[Live]") {
ServerSentEventParser parser;
const std::string stream = ": keepalive\n\nevent: file\ndata: {\"file_number\":3}\n\nevent: end\ndata: {}\n\n";
std::vector<ServerSentEvent> events;
// Fed in pieces, as the network delivers it.
for (size_t i = 0; i < stream.size(); i += 5) {
auto e = parser.Feed(stream.data() + i, std::min<size_t>(5, stream.size() - i));
events.insert(events.end(), e.begin(), e.end());
}
REQUIRE(events.size() == 2);
CHECK(events[0].event == "file");
CHECK(events[0].data == "{\"file_number\":3}");
CHECK(events[1].event == "end");
}
TEST_CASE("FileWriter_FileClosedCallback", "[Live]") {
RegisterHDF5Filter();
std::filesystem::remove_all(live_dir);
const auto x = LiveExperiment(live_dir + "/cb");
StartMessage start;
x.FillMessage(start);
std::vector<std::pair<uint64_t, uint64_t>> closed; // file number (from 1), images
{
FileWriter writer(start);
writer.FileClosedCallback([&](const HDF5DataFileStatistics &s) {
// Reported once the file is in place under its own name.
CHECK(std::filesystem::exists(s.filename));
closed.emplace_back(s.file_number, s.max_image_number + 1);
});
for (int i = 0; i < 7; i++)
writer.Write(Message(x, Image(x, i), i));
CHECK(closed.size() == 3); // the full files close as they fill
EndMessage end;
end.max_image_number = 7;
writer.WriteHDF5(end);
writer.Finalize();
}
REQUIRE(closed.size() == 4);
CHECK(closed[3] == std::make_pair<uint64_t, uint64_t>(4, 1)); // the partial last file, at END
std::filesystem::remove_all(live_dir);
}
TEST_CASE("TCPStreamPusher_ReportsClosedFiles", "[Live][TCP]") {
RegisterHDF5Filter();
Logger logger("TCPStreamPusher_ReportsClosedFiles");
std::filesystem::remove_all(live_dir);
std::filesystem::create_directories(live_dir);
const auto x = LiveExperiment(live_dir + "/tcp");
TCPStreamPusher pusher("tcp://127.0.0.1:*", 2);
std::vector<std::unique_ptr<TCPImagePuller>> pullers;
std::vector<std::unique_ptr<StreamWriter>> writers;
std::vector<std::future<StreamWriterOutput>> futures;
for (int i = 0; i < 2; i++) {
pullers.push_back(std::make_unique<TCPImagePuller>(pusher.GetAddress()[0], 8 * 1024 * 1024));
writers.push_back(std::make_unique<StreamWriter>(logger, *pullers.back()));
futures.push_back(std::async(std::launch::async, [w = writers.back().get()] { return w->Run(); }));
}
for (int i = 0; i < 100 && pusher.GetConnectedWriters() < 2; i++)
std::this_thread::sleep_for(std::chrono::milliseconds(50));
REQUIRE(pusher.GetConnectedWriters() == 2);
StartMessage start;
x.FillMessage(start);
start.write_master_file = true;
pusher.StartDataCollection(start);
LiveCollection *live = pusher.GetLiveCollection();
REQUIRE(live != nullptr);
REQUIRE(live->StartCBOR());
std::vector<uint8_t> buffer(16 * 1024 * 1024);
CBORStream2Serializer serializer(buffer.data(), buffer.size());
for (int i = 0; i < 7; i++) {
const auto compressed = Image(x, i);
serializer.SerializeImage(Message(x, compressed, i));
REQUIRE(pusher.SendImage(buffer.data(), serializer.GetBufferSize(), i));
}
EndMessage end;
end.max_image_number = 7;
REQUIRE(pusher.EndDataCollection(end));
const auto s = live->Get();
CHECK(s.ended);
CHECK(s.end_images == 7);
REQUIRE(s.files.size() == 4);
uint64_t images = 0;
for (const auto &f: s.files) {
CHECK(std::filesystem::exists(std::filesystem::path(live_dir) / std::filesystem::path(f.file).filename()));
CHECK(f.first_image == f.file_number * 2);
images += f.image_count;
}
CHECK(images == 7);
CHECK(std::filesystem::exists(live_dir + "/tcp_master.h5"));
for (auto &w: writers)
w->Cancel();
for (auto &f: futures)
REQUIRE_NOTHROW(f.get());
std::filesystem::remove_all(live_dir);
}
TEST_CASE("JFJochStartMessageReader", "[Live]") {
RegisterHDF5Filter();
std::filesystem::remove_all(live_dir);
const auto x = LiveExperiment(live_dir + "/sm");
StartMessage start;
x.FillMessage(start);
std::vector<std::vector<uint8_t>> images;
for (int i = 0; i < 7; i++)
images.push_back(Image(x, i));
JFJochStartMessageReader reader;
std::vector<std::pair<uint64_t, uint64_t>> closed;
{
FileWriter writer(start);
writer.FileClosedCallback([&](const HDF5DataFileStatistics &s) {
closed.emplace_back(s.file_number - 1, s.max_image_number + 1);
});
reader.Open(start, live_dir + "/sm_master.h5");
CHECK(reader.GetNumberOfImages() == 7);
// A read waits for its file.
auto read = std::async(std::launch::async, [&] { return reader.GetRawImage(2); });
for (int i = 0; i < 4; i++)
writer.Write(Message(x, images[i], i));
CHECK(read.wait_for(std::chrono::milliseconds(200)) == std::future_status::timeout);
for (const auto &[f, n]: closed)
reader.FileClosed(f, n);
auto raw = read.get();
REQUIRE(raw->image_buffer.size() == images[2].size());
CHECK(memcmp(raw->image_buffer.data(), images[2].data(), images[2].size()) == 0);
// Stopped after 5 images: the rest is never read.
writer.Write(Message(x, images[4], 4));
EndMessage end;
end.max_image_number = 5;
writer.WriteHDF5(end);
writer.Finalize();
}
for (const auto &[f, n]: closed)
reader.FileClosed(f, n);
reader.Ended();
JFJochReaderRawImage raw;
CHECK(reader.ReadRawImage(4, raw));
CHECK_FALSE(reader.ReadRawImage(5, raw));
CHECK_FALSE(reader.ReadRawImage(6, raw));
// The dataset is the one the master file describes.
JFJochHDF5Reader offline;
offline.ReadFile(live_dir + "/sm_master.h5");
const auto &a = reader.GetDataset()->experiment;
const auto &b = offline.GetDataset()->experiment;
CHECK(a.GetBeamX_pxl() == b.GetBeamX_pxl());
CHECK(a.GetBeamY_pxl() == b.GetBeamY_pxl());
CHECK(a.GetDetectorDistance_mm() == Catch::Approx(b.GetDetectorDistance_mm()));
CHECK(a.GetWavelength_A() == Catch::Approx(b.GetWavelength_A()));
CHECK(a.GetPixelSize_mm() == Catch::Approx(b.GetPixelSize_mm()));
CHECK(a.GetXPixelsNum() == b.GetXPixelsNum());
CHECK(a.GetYPixelsNum() == b.GetYPixelsNum());
REQUIRE(a.GetGoniometer());
REQUIRE(b.GetGoniometer());
CHECK(a.GetGoniometer()->GetStart_deg() == Catch::Approx(b.GetGoniometer()->GetStart_deg()));
CHECK(a.GetGoniometer()->GetIncrement_deg() == Catch::Approx(b.GetGoniometer()->GetIncrement_deg()));
CHECK(reader.GetDataset()->pixel_mask->GetMask() == offline.GetDataset()->pixel_mask->GetMask());
offline.Close();
// A broker that is gone fails a read that would wait.
JFJochStartMessageReader lost;
lost.Open(start, live_dir + "/sm_master.h5");
lost.Lost("gone");
CHECK_THROWS_AS(lost.ReadRawImage(0, raw), LiveCollectionLost);
reader.Close();
std::filesystem::remove_all(live_dir);
REQUIRE(H5Fget_obj_count(H5F_OBJ_ALL, H5F_OBJ_ALL) == 0);
}
TEST_CASE("BrokerFeed_StartFromEventStream", "[Live]") {
// The start message arrives as the first event, after a pause (a DECTRIS detector sends its own
// only after /start); the feed waits for it, then hands the reader the files and the end.
const auto x = LiveExperiment(live_dir + "/feed");
StartMessage start;
x.FillMessage(start);
std::vector<uint8_t> buffer(64 * 1024 * 1024);
CBORStream2Serializer serializer(buffer.data(), buffer.size());
serializer.SerializeSequenceStart(start);
const std::string cbor(buffer.begin(), buffer.begin() + serializer.GetBufferSize());
httplib::Server server;
server.Get("/live/events", [&](const httplib::Request &req, httplib::Response &res) {
REQUIRE(req.get_header_value("Authorization") == "Bearer secret");
res.set_chunked_content_provider("text/event-stream", [&, step = 0](size_t, httplib::DataSink &sink) mutable {
std::string out;
switch (step++) {
case 0:
std::this_thread::sleep_for(std::chrono::milliseconds(300));
out = ": keepalive\n\n";
break;
case 1:
out = "event: start\ndata: " + nlohmann::json{{"cbor", macaron::Encode(cbor)}}.dump() + "\n\n";
break;
case 2:
out = "event: file\ndata: {\"file_number\":0,\"image_count\":2}\n\n"
"event: end\ndata: {\"images\":2}\n\n";
break;
default:
sink.done();
return true;
}
return sink.write(out.data(), out.size());
});
});
const int port = server.bind_to_any_port("127.0.0.1");
std::thread server_thread([&] { server.listen_after_bind(); });
server.wait_until_ready();
{
BrokerFeed feed("http://127.0.0.1:" + std::to_string(port) + "/", "secret");
const StartMessage received = feed.WaitForStart();
CHECK(received.file_prefix == start.file_prefix);
CHECK(received.number_of_images == 7);
CHECK(received.beam_center_x == start.beam_center_x);
JFJochStartMessageReader reader;
reader.Open(received, live_dir + "/feed_master.h5");
feed.Attach(reader);
JFJochReaderRawImage raw;
// File 1 was never reported and the collection ended: no image there, and no wait.
CHECK_FALSE(reader.ReadRawImage(2, raw));
}
server.stop();
server_thread.join();
}