From 2981aedd3ed996462208fcd2aa1eb7ed0c0ce72f Mon Sep 17 00:00:00 2001 From: Filip Leonarski Date: Thu, 8 Oct 2026 17:03:09 +0200 Subject: [PATCH] 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) Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi --- broker/JFJochBrokerHttp.cpp | 34 ++----- broker/JFJochBrokerHttp.h | 5 +- broker/JFJochStateMachine.cpp | 15 +++ broker/jfjoch_api.yaml | 41 ++------ docs/JFJOCH_BROKER.md | 26 ++--- docs/SECURITY.md | 2 +- image_pusher/LiveCollection.cpp | 30 +++++- image_pusher/LiveCollection.h | 14 ++- rugnux/BrokerFeed.cpp | 165 +++++++++++++++++++------------- rugnux/BrokerFeed.h | 35 +++++-- rugnux/rugnux_cli.cpp | 6 +- tests/BrokerHttpAuthTest.cpp | 3 +- tests/LiveCollectionTest.cpp | 103 +++++++++++++++++++- 13 files changed, 315 insertions(+), 164 deletions(-) diff --git a/broker/JFJochBrokerHttp.cpp b/broker/JFJochBrokerHttp.cpp index 67fea4754..b722a6600 100644 --- a/broker/JFJochBrokerHttp.cpp +++ b/broker/JFJochBrokerHttp.cpp @@ -16,6 +16,7 @@ #include "../preview/JFJochTIFF.h" #include "OpenAPIConvert.h" #include "../image_pusher/LiveCollection.h" +#include #include "gen/model/Error_message.h" using namespace org::openapitools::server::model; @@ -151,7 +152,6 @@ static const std::set PROTECTED_PATHS{ "/image_buffer/image.tiff", "/preview/plot", "/preview/plot.bin", - "/live/start.cbor", "/live/events" }; @@ -331,7 +331,6 @@ void JFJochBrokerHttp::register_routes(httplib::Server &server) { }); server.Get("/image_buffer/start.cbor", bind_noarg(&JFJochBrokerHttp::image_buffer_start_cbor_get)); - server.Get("/live/start.cbor", bind_noarg(&JFJochBrokerHttp::live_start_cbor_get)); server.Get("/live/events", bind_noarg(&JFJochBrokerHttp::live_events_get)); server.Get("/image_buffer/status", bind_noarg(&JFJochBrokerHttp::image_buffer_status_get)); server.Get("/image_pusher/status", bind_noarg(&JFJochBrokerHttp::image_pusher_status_get)); @@ -794,22 +793,6 @@ void JFJochBrokerHttp::image_buffer_start_cbor_get(httplib::Response &response) response.status = 404; } -void JFJochBrokerHttp::live_start_cbor_get(httplib::Response &response) { - LiveCollection *live = services.GetLiveCollection(); - if (live == nullptr) { - response.status = 409; - response.set_content("Following a collection needs the TCP or the HDF5 image pusher", "text/plain"); - return; - } - const auto cbor = live->StartCBOR(); - if (!cbor) { - response.status = 404; - response.set_content("No collection has started", "text/plain"); - return; - } - response.set_content(std::string(cbor->begin(), cbor->end()), "application/cbor"); -} - namespace { std::string ServerSentEvent(const std::string &event, const nlohmann::json &data) { return "event: " + event + "\ndata: " + data.dump() + "\n\n"; @@ -853,15 +836,18 @@ void JFJochBrokerHttp::live_events_get(httplib::Response &response) { sink.done(); return true; } - const auto s = live->Wait(generation, cursor->files_sent, - cursor->collection_sent ? std::chrono::seconds(10) - : std::chrono::seconds(0)); + const auto s = live->Wait(generation, cursor->collection_sent, cursor->files_sent, + std::chrono::seconds(10)); std::string out; - if (!cursor->collection_sent) { - out += ServerSentEvent("collection", { + // The start message first, once it is known: a collection is announced at /start, and a + // DECTRIS detector sends its start message only afterwards. + if (!cursor->collection_sent && s.started && s.generation == generation) { + const auto cbor = live->StartCBOR().value_or(std::vector{}); + out += ServerSentEvent("start", { {"run_number", s.run_number}, {"run_name", s.run_name}, {"file_prefix", s.file_prefix}, {"number_of_images", s.number_of_images}, - {"images_per_file", s.images_per_file}}); + {"images_per_file", s.images_per_file}, + {"cbor", macaron::Encode(std::string(cbor.begin(), cbor.end()))}}); cursor->collection_sent = true; } if (s.generation != generation) { diff --git a/broker/JFJochBrokerHttp.h b/broker/JFJochBrokerHttp.h index fd8bced31..8d42a4593 100644 --- a/broker/JFJochBrokerHttp.h +++ b/broker/JFJochBrokerHttp.h @@ -85,10 +85,9 @@ class JFJochBrokerHttp { void register_routes(httplib::Server &server); - // A reader following the current collection while it is written: its start message, and an - // event stream of the data files as they are closed (docs/JFJOCH_BROKER.md, "Following a + // A reader following the current collection while it is written: an event stream of its start + // message and of the data files as they are closed (docs/JFJOCH_BROKER.md, "Following a // collection"). - void live_start_cbor_get(httplib::Response &response); void live_events_get(httplib::Response &response); void cancel_post(httplib::Response &response); diff --git a/broker/JFJochStateMachine.cpp b/broker/JFJochStateMachine.cpp index 398fdf7ee..1b93ad2a7 100644 --- a/broker/JFJochStateMachine.cpp +++ b/broker/JFJochStateMachine.cpp @@ -4,6 +4,7 @@ #include #include "JFJochStateMachine.h" +#include "../image_pusher/LiveCollection.h" #include "../preview/JFJochTIFF.h" #include "../common/CUDAWrapper.h" #include "../common/GitInfo.h" @@ -426,6 +427,10 @@ void JFJochStateMachine::Start(const DatasetSettings &settings, bool async, std: experiment.IncrementRunNumber(); SetState(JFJochState::Busy, "Preparing measurement", BrokerStatus::MessageSeverity::Info); + // From here on a reader can follow the collection (/live/*), even before the image pusher has + // its start message. + if (auto *live = services.GetLiveCollection()) + live->Expect(); measurement = std::async(std::launch::async, &JFJochStateMachine::MeasurementThread, this); if (!async) { c.wait(ul, [&]() { return state != JFJochState::Busy; }); @@ -465,6 +470,16 @@ PixelMaskStatistics JFJochStateMachine::GetPixelMaskStatistics() const { } void JFJochStateMachine::MeasurementThread() { + // However the measurement ends, a reader following it (/live/events) learns that it did - also + // when it ended before the image pusher started or ended it. + struct LiveEnd { + LiveCollection *live; + ~LiveEnd() { + if (live) + live->EndIfOpen("the measurement ended without the writers finishing the collection"); + } + } live_end{services.GetLiveCollection()}; + try { // Before anything is started: ask the writer whether it could write this run. A run that // would be refused for a file already there, or a directory that cannot be made, is diff --git a/broker/jfjoch_api.yaml b/broker/jfjoch_api.yaml index 64e93a8b5..0370f3029 100644 --- a/broker/jfjoch_api.yaml +++ b/broker/jfjoch_api.yaml @@ -3867,45 +3867,22 @@ paths: application/json: schema: $ref: '#/components/schemas/error_message' - /live/start.cbor: - get: - summary: Start message of the collection, for a reader following it - description: | - The CBOR start message the writer of the master file received for the current collection. A - reader (rugnux) builds the dataset from it and then reads each data file once - /live/events reports it closed. Needs the TCP or the HDF5 image pusher. - security: - - bearerAuth: [] - responses: - "200": - description: Start message - content: - application/cbor: - schema: - type: string - format: binary - "401": - description: The current dataset is protected and no valid bearer token was presented. - content: - text/plain: - schema: - type: string - "404": - description: No collection has started - "409": - description: The image pusher does not report closed data files /live/events: get: summary: Event stream of the collection, for a reader following it description: | - Server-sent events (text/event-stream) for the current collection: first its whole history, - then every change, until it ends or a new collection begins. Only one reader at a time. - Events, each with a JSON `data` line: - - `collection` {run_number, run_name, file_prefix, number_of_images, images_per_file} - first; + Server-sent events (text/event-stream) for the current collection, from the moment /start has + accepted it: first its whole history, then every change, until it ends or a new collection + begins. Only one reader at a time. Events, each with a JSON `data` line: + - `start` {run_number, run_name, file_prefix, number_of_images, images_per_file, cbor} - first, + once the start message is known: at once for a JUNGFRAU, when the detector sends its own for + a DECTRIS. `cbor` is the start message the writer of the master file received, CBOR in base64; + a reader builds the dataset from it; - `file` {file_number, file, first_image, image_count} - a data file closed and renamed into place; `file` is relative to the collection's root, as the writer named it; - `end` {images, error?} - every writer acknowledged the end: the data files and the master - file are in place; `images` is the number of images collected. The stream then closes; + file are in place; `images` is the number of images collected. Also sent, with `error`, when + the measurement ends any other way - before the start message, for one. The stream then closes; - `superseded` {run_number} - a new collection began. The stream then closes. A comment line (`: keepalive`) every 10 s when nothing happens. security: diff --git a/docs/JFJOCH_BROKER.md b/docs/JFJOCH_BROKER.md index 2253a466c..739df1d68 100644 --- a/docs/JFJOCH_BROKER.md +++ b/docs/JFJOCH_BROKER.md @@ -153,22 +153,22 @@ Example configuration (not every section is shown): ## Following a collection A reader can process a collection while it is being written - rugnux does, given the path the master -file will have and the broker's address (`rugnux _master.h5 http://:/`). The -images still travel as files; the broker says which files exist. Two routes serve it, protected by the -dataset's bearer tokens like the other dataset routes: - -* `GET /live/start.cbor` - the CBOR start message the writer of the master file received. A reader - builds the dataset from it. -* `GET /live/events` - server-sent events for the current collection: its history first, then every - change. `collection` (run number, file prefix, image count, images per file) comes first; `file` - (file number, file name relative to the collection's root, first image, image count) for each data - file a writer has closed and renamed into place; `end` (images collected, and an error if any) once - every writer has acknowledged the end, so the data files and the master file are in place. A - comment line every 10 s keeps the connection alive. One reader at a time; a second one gets 409. +file will have and the broker's address (`rugnux _master.h5 http://:/`), started +any time after `/start` was accepted. The images still travel as files; the broker says which files +exist. `GET /live/events`, protected by the dataset's bearer tokens like the other dataset routes, +streams server-sent events for the current collection from the moment `/start` accepted it: its +history first, then every change. `start` comes first, once the start message is known - at once +for a JUNGFRAU, when the detector sends its own for a DECTRIS - with that start message (CBOR, in +base64); a reader builds the dataset from it. Then `file` (file number, file name relative to the +collection's root, first image, image count) for each data file a writer has closed and renamed into +place, and `end` (images collected, and an error if any) once every writer has acknowledged the end, +so the data files and the master file are in place - or, with an error, when the measurement ended +any other way. A comment line every 10 s keeps the connection alive. One reader at a time; a second +one gets 409. The writers report each closed data file to the broker (TCP protocol version 5, an acknowledgement for `FILE_CLOSED`); the in-process HDF5 writer reports to it directly. The ZeroMQ writer has no way -back to the broker, so with it the routes answer 409. +back to the broker, so with it the route answers 409. ## Setting up a local test for Jungfraujoch For development, it is possible to set up a local installation of Jungfraujoch. diff --git a/docs/SECURITY.md b/docs/SECURITY.md index 84639d237..4893b63a2 100644 --- a/docs/SECURITY.md +++ b/docs/SECURITY.md @@ -109,7 +109,7 @@ served under the new tokens, and the new run's name never under the old ones. A | `/result/scan` (dataset name, cell of a grid scan / rotation) | 401 without a token | | `/image_buffer/start.cbor`, `/image_buffer/image.cbor`, `/image_buffer/image.jpeg`, `/image_buffer/image.tiff` (the images and the start message) | 401 without a token | | `/preview/plot`, `/preview/plot.bin` (per-image plots, unit cell) | 401 without a token | -| `/live/start.cbor`, `/live/events` (the start message and the data files of the collection being written) | 401 without a token | +| `/live/events` (the start message and the data files of the collection being written) | 401 without a token | | `/statistics` (the aggregate the web UI polls) | 200, but the `measurement` block is omitted | | everything else (`/status`, `/config/*`, `/start`, `/cancel`, masks, pedestal, ...) | open | diff --git a/image_pusher/LiveCollection.cpp b/image_pusher/LiveCollection.cpp index a64c26c26..d40f5ac56 100644 --- a/image_pusher/LiveCollection.cpp +++ b/image_pusher/LiveCollection.cpp @@ -5,12 +5,27 @@ #include "../writer/HDF5NXmx.h" -void LiveCollection::Begin(const StartMessage &start, std::vector cbor) { +void LiveCollection::Expect() { { std::unique_lock ul(m); const uint64_t generation = state.generation + 1; state = Snapshot{}; state.generation = generation; + start_cbor.clear(); + } + cv.notify_all(); +} + +void LiveCollection::Begin(const StartMessage &start, std::vector cbor) { + { + std::unique_lock ul(m); + // The collection Expect announced, or - when nothing announced it - a new one. + if (state.generation == 0 || state.started || state.ended) { + const uint64_t generation = state.generation + 1; + state = Snapshot{}; + state.generation = generation; + } + state.started = true; state.run_number = start.run_number; state.run_name = start.run_name; state.file_prefix = start.file_prefix; @@ -25,7 +40,7 @@ void LiveCollection::Begin(const StartMessage &start, std::vector cbor) void LiveCollection::FileClosed(uint64_t file_number, uint64_t image_count) { { std::unique_lock ul(m); - if (state.generation == 0 || state.ended) + if (!state.started || state.ended) return; state.files.push_back(DataFile{ .file_number = file_number, @@ -49,18 +64,23 @@ void LiveCollection::End(uint64_t images, const std::string &error) { cv.notify_all(); } +void LiveCollection::EndIfOpen(const std::string &error) { + End(0, error); +} + std::optional> LiveCollection::StartCBOR() const { std::unique_lock ul(m); - if (state.generation == 0) + if (!state.started) return {}; return start_cbor; } -LiveCollection::Snapshot LiveCollection::Wait(uint64_t generation, size_t files_seen, +LiveCollection::Snapshot LiveCollection::Wait(uint64_t generation, bool started_seen, size_t files_seen, std::chrono::milliseconds wait) { std::unique_lock ul(m); cv.wait_for(ul, wait, [&] { - return state.generation != generation || state.ended || state.files.size() > files_seen; + return state.generation != generation || state.ended || state.files.size() > files_seen + || (state.started && !started_seen); }); return state; } diff --git a/image_pusher/LiveCollection.h b/image_pusher/LiveCollection.h index 157ed79cd..d35cfbf67 100644 --- a/image_pusher/LiveCollection.h +++ b/image_pusher/LiveCollection.h @@ -28,6 +28,7 @@ public: struct Snapshot { uint64_t generation = 0; // 0: no collection since the broker started + bool started = false; // the start message is known uint64_t run_number = 0; std::string run_name; std::string file_prefix; @@ -39,17 +40,24 @@ public: std::string end_error; }; + // A collection was accepted (/start); its start message comes with Begin, once the image pusher + // starts the writers - for a DECTRIS detector only when the detector's own start message arrives. + // A reader following it in between waits for it. + void Expect(); // start_cbor: the start message as the writer of the master file received it. void Begin(const StartMessage &start, std::vector start_cbor); void FileClosed(uint64_t file_number, uint64_t image_count); // After every writer has acknowledged END: the data files and the master are in place. void End(uint64_t images, const std::string &error = ""); + // The measurement is over: ends a collection the image pusher did not end, with this error. + void EndIfOpen(const std::string &error); + // The start message of the current collection; nothing before it is known. std::optional> StartCBOR() const; - // Blocks until the collection has more than `files_seen` files, has ended, a new collection has - // begun, or `wait` passes; returns the state then. - Snapshot Wait(uint64_t generation, size_t files_seen, std::chrono::milliseconds wait); + // Blocks until the start message arrives (when not `started_seen`), the collection has more than + // `files_seen` files, has ended, a new collection has begun, or `wait` passes; returns the state then. + Snapshot Wait(uint64_t generation, bool started_seen, size_t files_seen, std::chrono::milliseconds wait); Snapshot Get() const; // The single follower slot; false when another follower holds it. diff --git a/rugnux/BrokerFeed.cpp b/rugnux/BrokerFeed.cpp index c5b02dec3..2a1c346d5 100644 --- a/rugnux/BrokerFeed.cpp +++ b/rugnux/BrokerFeed.cpp @@ -3,6 +3,7 @@ #include "BrokerFeed.h" +#include #include #include @@ -41,12 +42,15 @@ namespace { return {{"Authorization", "Bearer " + token}}; } - std::string Refused(const httplib::Result &res, const std::string &what) { - if (!res) - return what + ": no answer from the broker (" + httplib::to_string(res.error()) + ")"; - if (res->status == 401) - return what + ": the broker refused the token in JUNGFRAUJOCH_HTTP_TOKEN"; - return what + ": the broker answered " + std::to_string(res->status) + " " + res->body; + std::vector DecodeBase64(const std::string &s) { + size_t n = s.size() / 4 * 3; + if (!s.empty() && s.back() == '=') + n--; + if (s.size() > 1 && s[s.size() - 2] == '=') + n--; + std::vector out(n); + macaron::Decode(s, out); + return out; } } @@ -66,6 +70,7 @@ BrokerFeed::BrokerFeed(const std::string &url, std::string in_token) : token(std prefix = url.substr(path); while (prefix.ends_with('/')) prefix.pop_back(); + thread = std::thread([this] { Run(); }); } BrokerFeed::~BrokerFeed() { @@ -74,69 +79,93 @@ BrokerFeed::~BrokerFeed() { thread.join(); } -StartMessage BrokerFeed::FetchStart() { - httplib::Client client(host); - auto res = client.Get(prefix + "/live/start.cbor", Authorization(token)); - if (!res || res->status != 200) - throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, - Refused(res, "No collection to follow at " + host)); - auto msg = CBORStream2Deserialize(res->body); - if (!msg || !msg->start_message) - throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, - "The broker at " + host + " sent no start message"); - return msg->start_message.value(); +StartMessage BrokerFeed::WaitForStart() { + std::unique_lock ul(m); + cv.wait(ul, [&] { return start.has_value() || done; }); + if (start) + return start.value(); + throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, + lost.empty() ? "The broker at " + host + " sent no start message" : lost); } -void BrokerFeed::Follow(JFJochStartMessageReader &reader, uint64_t run_number, const std::string &file_prefix) { - thread = std::thread([this, &reader, run_number, file_prefix] { - Logger logger("BrokerFeed"); - httplib::Client client(host); - // The broker writes at least a keepalive every 10 s; three missed ones mean it is gone. - client.set_read_timeout(std::chrono::seconds(30)); - ServerSentEventParser parser; - bool ended = false; - std::string lost; - size_t files = 0; - int status = 0; - auto res = client.Get(prefix + "/live/events", Authorization(token), - [&](const httplib::Response &response) { - status = response.status; - return status == 200; - }, - [&](const char *data, size_t size) { - for (const auto &e: parser.Feed(data, size)) { - const auto j = nlohmann::json::parse(e.data); - if (e.event == "collection") { - if (j.at("run_number").get() != run_number - || j.at("file_prefix").get() != file_prefix) - lost = "the broker is now on another collection (" + j.at("file_prefix").get() + ")"; - } else if (e.event == "file") { - reader.FileClosed(j.at("file_number").get(), j.at("image_count").get()); - files++; - } else if (e.event == "end") { - if (j.contains("error")) - logger.Warning("The collection ended with an error: {}", j.at("error").get()); - logger.Info("The collection ended with {} images in {} data files", - j.at("images").get(), files); - ended = true; - } else if (e.event == "superseded") { - lost = "a new collection began before this one ended"; - } - if (!lost.empty()) - return false; - } - return !stop; - }); - if (ended) - reader.Ended(); - else if (!lost.empty()) - reader.Lost(lost); +void BrokerFeed::Attach(JFJochStartMessageReader &in_reader) { + std::unique_lock ul(m); + reader = &in_reader; + for (const auto &[file_number, images]: files) + reader->FileClosed(file_number, images); + if (ended) + reader->Ended(); + else if (done) + reader->Lost(lost); +} + +void BrokerFeed::Run() { + Logger logger("BrokerFeed"); + httplib::Client client(host); + // The broker writes at least a keepalive every 10 s; three missed ones mean it is gone. + client.set_read_timeout(std::chrono::seconds(30)); + ServerSentEventParser parser; + int status = 0; + std::string error; + client.Get(prefix + "/live/events", Authorization(token), + [&](const httplib::Response &response) { + status = response.status; + return status == 200; + }, + [&](const char *data, size_t size) { + for (const auto &e: parser.Feed(data, size)) { + const auto j = nlohmann::json::parse(e.data); + std::unique_lock ul(m); + if (e.event == "start") { + auto msg = CBORStream2Deserialize(DecodeBase64(j.at("cbor").get())); + if (msg && msg->start_message) + start = msg->start_message.value(); + else + error = "the broker sent no readable start message"; + } else if (e.event == "file") { + files.emplace_back(j.at("file_number").get(), j.at("image_count").get()); + if (reader) + reader->FileClosed(files.back().first, files.back().second); + } else if (e.event == "end") { + if (j.contains("error")) + logger.Warning("The collection ended with an error: {}", j.at("error").get()); + if (!start) + error = "The collection at " + host + " ended before it started" + + (j.contains("error") ? ": " + j.at("error").get() : ""); + else { + logger.Info("The collection ended with {} images in {} data files", + j.at("images").get(), files.size()); + ended = true; + } + } else if (e.event == "superseded") { + error = "a new collection began before this one ended"; + } + cv.notify_all(); + if (!error.empty()) + return false; + } + return !stop; + }); + + std::unique_lock ul(m); + if (!ended) { + if (!error.empty()) + lost = error; else if (status == 401) - reader.Lost("Following the collection at " + host + ": the broker refused the token in JUNGFRAUJOCH_HTTP_TOKEN"); - else if (status != 200 && status != 0) - reader.Lost("Following the collection at " + host + ": the broker answered " + std::to_string(status) - + (status == 409 ? " (another reader is following it)" : "")); - else if (!stop) - reader.Lost("Lost the broker at " + host + " before the collection ended"); - }); + lost = "Following the collection at " + host + ": the broker refused the token in JUNGFRAUJOCH_HTTP_TOKEN"; + else if (status == 404) + lost = "No collection to follow at " + host; + else if (status == 409) + lost = "Following the collection at " + host + ": another reader is following it, or its writer " + "does not report closed data files"; + else if (status != 200) + lost = "No answer from the broker at " + host; + else + lost = "Lost the broker at " + host + " before the collection ended"; + if (reader) + reader->Lost(lost); + } else if (reader) + reader->Ended(); + done = true; + cv.notify_all(); } diff --git a/rugnux/BrokerFeed.h b/rugnux/BrokerFeed.h index 0e21b859e..4eeeedba7 100644 --- a/rugnux/BrokerFeed.h +++ b/rugnux/BrokerFeed.h @@ -4,6 +4,9 @@ #pragma once #include +#include +#include +#include #include #include #include @@ -25,30 +28,44 @@ public: std::vector Feed(const char *bytes, size_t size); }; -// rugnux's side of following a collection while it is written: the start message from the broker -// (GET /live/start.cbor), then its event stream (GET /live/events) handed to a JFJochStartMessageReader. -// The token goes in the Authorization header and nowhere else. +// rugnux's side of following a collection while it is written: the broker's event stream +// (GET /live/events), read in a thread of its own from construction on. Its first event is the +// collection's start message; then every closed data file, handed to a JFJochStartMessageReader; then +// the end. The token goes in the Authorization header and nowhere else. class BrokerFeed { public: // url: http://host:port[/prefix][/] BrokerFeed(const std::string &url, std::string token); ~BrokerFeed(); - // Throws with the broker's answer when there is no collection to follow. - StartMessage FetchStart(); + // Waits for the start message - at once on a JUNGFRAU, when the detector sends its own on a + // DECTRIS. Throws with the reason when the stream ends or fails first (no collection, a refused + // token, a collection that ended before it started). + StartMessage WaitForStart(); - // Reads the event stream in a thread of its own until the collection ends, reporting every closed - // data file to `reader`. The reader is told the broker is lost when the stream breaks before the - // end, or when it describes another collection than run_number / file_prefix. - void Follow(JFJochStartMessageReader &reader, uint64_t run_number, const std::string &file_prefix); + // Hands `reader` every closed data file, those already reported and those to come, and tells it + // when the collection ends or the broker is lost. + void Attach(JFJochStartMessageReader &reader); // scheme://host:port, for messages - never the token. [[nodiscard]] const std::string &Host() const { return host; } private: + void Run(); + std::string host; // http://host:port std::string prefix; // path before /live, without a trailing slash std::string token; std::atomic stop{false}; + + std::mutex m; + std::condition_variable cv; + std::optional start; + std::vector> files; // file number, images + bool ended = false; + bool done = false; + std::string lost; + JFJochStartMessageReader *reader = nullptr; + std::thread thread; }; diff --git a/rugnux/rugnux_cli.cpp b/rugnux/rugnux_cli.cpp index 75c96dc7c..9e74ef787 100644 --- a/rugnux/rugnux_cli.cpp +++ b/rugnux/rugnux_cli.cpp @@ -1621,7 +1621,9 @@ static int RunRugnux(int argc, char **argv) { exit(EXIT_FAILURE); } broker_feed = std::make_unique(broker_url, token); - const StartMessage start = broker_feed->FetchStart(); + // The broker's event stream starts with the start message: at once on a JUNGFRAU, when the + // detector sends its own on a DECTRIS. + const StartMessage start = broker_feed->WaitForStart(); // The start message names the files relative to the collection's root; the path given here // has to be where its master file will be. const std::string master = std::filesystem::path(HDF5Metadata::MasterFileName(start)).generic_string(); @@ -1635,7 +1637,7 @@ static int RunRugnux(int argc, char **argv) { exit(EXIT_FAILURE); } live_reader.Open(start, input_file); - broker_feed->Follow(live_reader, start.run_number, start.file_prefix); + broker_feed->Attach(live_reader); reader_ptr = &live_reader; logger.Info("Following collection {} (run {}) from the broker at {}: {} images in files of {}", master, start.run_number, broker_feed->Host(), start.number_of_images, start.images_per_file); diff --git a/tests/BrokerHttpAuthTest.cpp b/tests/BrokerHttpAuthTest.cpp index 15489c049..cdcafff52 100644 --- a/tests/BrokerHttpAuthTest.cpp +++ b/tests/BrokerHttpAuthTest.cpp @@ -69,7 +69,7 @@ TEST_CASE("JFJochBrokerHttp_BearerTokens", "[broker]") { for (const char *path: {"/statistics/data_collection", "/result/scan", "/image_buffer/start.cbor", "/image_buffer/image.cbor?id=0", "/image_buffer/image.tiff?id=0", "/image_buffer/image.jpeg?id=0", "/preview/plot?type=bkg_estimate", - "/preview/plot.bin?type=bkg_estimate", "/live/start.cbor", "/live/events"}) { + "/preview/plot.bin?type=bkg_estimate", "/live/events"}) { INFO(path); auto res = client.Get(path); REQUIRE(res); @@ -100,7 +100,6 @@ TEST_CASE("JFJochBrokerHttp_BearerTokens", "[broker]") { } SECTION("with no receiver there is no collection to follow") { - REQUIRE(client.Get("/live/start.cbor", good)->status == 409); REQUIRE(client.Get("/live/events", good)->status == 409); } diff --git a/tests/LiveCollectionTest.cpp b/tests/LiveCollectionTest.cpp index 47f4210c4..d481eef6c 100644 --- a/tests/LiveCollectionTest.cpp +++ b/tests/LiveCollectionTest.cpp @@ -2,6 +2,8 @@ // SPDX-License-Identifier: GPL-3.0-only #include +#include +#include #include #include @@ -64,13 +66,13 @@ TEST_CASE("LiveCollection_Lifecycle", "[Live]") { live.Begin(start, {1, 2, 3}); CHECK(live.StartCBOR() == std::vector{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{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 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(); +}