So that rugnux can process a collection while it is written, the broker reports each data file once
a writer has closed it and renamed it into place:
- Writer: FileWriter takes a callback, called for every closed data file. StreamWriter uses it on TCP
to send an acknowledgement with ack_for = FILE_CLOSED (new TCPFrameType 10): file number and the
images the file holds. TCP protocol version 5. Files closed at END are reported before the END
acknowledgement.
- LiveCollection (image_pusher): the start message the writer of the master file received, the
closed files, the end. Fed by the TCP pusher (FILE_CLOSED acknowledgements; end once every writer
acknowledged END) and by the in-process HDF5 pusher (the FileWriter callback). The ZeroMQ pusher
has no back channel and offers none.
- Two routes, protected by the dataset's bearer tokens like the other dataset routes:
- GET /live/start.cbor: that start message;
- GET /live/events: server-sent events (collection, file, end, superseded) with the history first
and a keepalive every 10 s. One follower at a time; a second gets 409.
Documented in jfjoch_api.yaml; they add no schema, so the generated C++ model is unchanged.
- The HTTP server's thread pool is set explicitly (16, up to 64): an event stream holds a thread for
the whole collection.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SVmAWnzCmRKAXVUCdc4iNi
136 lines
6.4 KiB
C++
136 lines
6.4 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 <nlohmann/json.hpp>
|
|
#include <chrono>
|
|
#include <thread>
|
|
|
|
#include "../broker/JFJochBrokerHttp.h"
|
|
|
|
// The tokens given to /start guard the endpoints that expose the dataset, and nothing else, for as
|
|
// long as that dataset is the current one. Runs the real HTTP layer against a broker with no
|
|
// receiver: /start is then accepted and finishes at once with nothing to acquire.
|
|
TEST_CASE("JFJochBrokerHttp_BearerTokens", "[broker]") {
|
|
DiffractionExperiment experiment;
|
|
SpotFindingSettings spot_finding;
|
|
JFJochBrokerHttp broker(experiment, spot_finding);
|
|
broker.AddDetectorSetup(DetJF4M());
|
|
|
|
httplib::Server server;
|
|
broker.attach(server);
|
|
const int port = server.bind_to_any_port("127.0.0.1");
|
|
REQUIRE(port > 0);
|
|
// Stopped and joined however the test ends, so a failed REQUIRE does not abort on a joinable thread.
|
|
struct ServerThread {
|
|
httplib::Server &server;
|
|
std::thread thread;
|
|
explicit ServerThread(httplib::Server &s) : server(s), thread([&s] { s.listen_after_bind(); }) {}
|
|
~ServerThread() { server.stop(); thread.join(); }
|
|
} server_thread(server);
|
|
server.wait_until_ready();
|
|
|
|
httplib::Client client("127.0.0.1", port);
|
|
const httplib::Headers good{{"Authorization", "Bearer secret-1"}};
|
|
const httplib::Headers other{{"Authorization", "Bearer secret-2"}};
|
|
const httplib::Headers wrong{{"Authorization", "Bearer nope"}};
|
|
|
|
auto start = [&](const nlohmann::json &extra) {
|
|
nlohmann::json body = {{"beam_x_pxl", 1000}, {"beam_y_pxl", 1000},
|
|
{"detector_distance_mm", 100}, {"incident_energy_keV", 12.4},
|
|
{"images_per_trigger", 1}, {"ntrigger", 1},
|
|
{"file_prefix", "protected_run"}};
|
|
body.update(extra);
|
|
auto res = client.Post("/start", body.dump(), "application/json");
|
|
REQUIRE(res);
|
|
REQUIRE(res->status == 200);
|
|
// With no receiver the measurement thread idles for 30 s before the run ends; wait it out.
|
|
for (int i = 0; i < 120; i++) {
|
|
auto status = client.Get("/status");
|
|
REQUIRE(status);
|
|
if (nlohmann::json::parse(status->body).at("state") == "Idle")
|
|
return;
|
|
std::this_thread::sleep_for(std::chrono::milliseconds(500));
|
|
}
|
|
FAIL("broker did not return to Idle after the start");
|
|
};
|
|
|
|
REQUIRE(client.Post("/initialize")->status == 200);
|
|
|
|
// Nothing started yet: open, whatever the header says.
|
|
REQUIRE(client.Get("/statistics/data_collection")->status == 200);
|
|
REQUIRE(client.Get("/statistics/data_collection", wrong)->status == 200);
|
|
REQUIRE(client.Get("/statistics")->status == 200);
|
|
|
|
start({{"tokens", {"secret-1", "secret-2"}}});
|
|
|
|
SECTION("protected endpoints refuse without a token and say nothing about the dataset") {
|
|
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"}) {
|
|
INFO(path);
|
|
auto res = client.Get(path);
|
|
REQUIRE(res);
|
|
REQUIRE(res->status == 401);
|
|
REQUIRE(res->get_header_value("WWW-Authenticate") == "Bearer");
|
|
REQUIRE(res->body.find("protected_run") == std::string::npos);
|
|
REQUIRE(client.Get(path, wrong)->status == 401);
|
|
}
|
|
}
|
|
|
|
SECTION("either token opens them") {
|
|
auto res = client.Get("/statistics/data_collection", good);
|
|
REQUIRE(res->status == 200);
|
|
REQUIRE(res->body.find("protected_run") != std::string::npos);
|
|
REQUIRE(client.Get("/statistics/data_collection", other)->status == 200);
|
|
REQUIRE(client.Get("/preview/plot?type=bkg_estimate", good)->status == 200);
|
|
}
|
|
|
|
SECTION("the aggregate statistics drop the measurement block instead of refusing") {
|
|
auto res = client.Get("/statistics");
|
|
REQUIRE(res->status == 200);
|
|
REQUIRE_FALSE(nlohmann::json::parse(res->body).contains("measurement"));
|
|
REQUIRE(res->body.find("protected_run") == std::string::npos);
|
|
|
|
auto with_token = client.Get("/statistics", good);
|
|
REQUIRE(with_token->status == 200);
|
|
REQUIRE(nlohmann::json::parse(with_token->body).contains("measurement"));
|
|
}
|
|
|
|
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);
|
|
}
|
|
|
|
SECTION("the control plane and the broker status stay open") {
|
|
REQUIRE(client.Get("/status")->status == 200);
|
|
REQUIRE(client.Get("/image_buffer/status")->status == 200);
|
|
REQUIRE(client.Get("/config/detector")->status == 200);
|
|
}
|
|
|
|
SECTION("a refused start keeps the current tokens, an accepted one replaces them") {
|
|
// Not idle: /start is refused, the dataset stays protected by its own tokens.
|
|
REQUIRE(client.Post("/deactivate")->status == 200);
|
|
nlohmann::json body = {{"beam_x_pxl", 1000}, {"beam_y_pxl", 1000},
|
|
{"detector_distance_mm", 100}, {"incident_energy_keV", 12.4}};
|
|
REQUIRE(client.Post("/start", body.dump(), "application/json")->status == 500);
|
|
REQUIRE(client.Post("/initialize")->status == 200);
|
|
REQUIRE(client.Get("/statistics/data_collection")->status == 401);
|
|
REQUIRE(client.Get("/statistics/data_collection", good)->status == 200);
|
|
|
|
// The next start without tokens lifts the protection.
|
|
start(nlohmann::json::object());
|
|
REQUIRE(client.Get("/statistics/data_collection")->status == 200);
|
|
REQUIRE(client.Get("/statistics/data_collection", wrong)->status == 200);
|
|
|
|
// And with new tokens, the old ones no longer open it.
|
|
start({{"tokens", {"secret-3"}}});
|
|
const httplib::Headers newer{{"Authorization", "Bearer secret-3"}};
|
|
REQUIRE(client.Get("/statistics/data_collection", good)->status == 401);
|
|
REQUIRE(client.Get("/statistics/data_collection", newer)->status == 200);
|
|
}
|
|
|
|
}
|