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
This commit is contained in:
2026-10-08 17:03:09 +02:00
co-authored by Claude Opus 5.5
parent 72c7409327
commit 2981aedd3e
13 changed files with 315 additions and 164 deletions
+101 -2
View File
@@ -2,6 +2,8 @@
// SPDX-License-Identifier: GPL-3.0-only
#include <catch2/catch_all.hpp>
#include <httplib.h>
#include <base64/Base64.h>
#include <filesystem>
#include <future>
@@ -64,13 +66,13 @@ TEST_CASE("LiveCollection_Lifecycle", "[Live]") {
live.Begin(start, {1, 2, 3});
CHECK(live.StartCBOR() == std::vector<uint8_t>{1, 2, 3});
auto s = live.Wait(1, 0, std::chrono::milliseconds(0));
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, 0, std::chrono::seconds(10)); });
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();
@@ -102,6 +104,48 @@ TEST_CASE("LiveCollection_Lifecycle", "[Live]") {
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";
@@ -278,3 +322,58 @@ TEST_CASE("JFJochStartMessageReader", "[Live]") {
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();
}