Files
Jungfraujoch/reader/JFJochStartMessageReader.h
leonarski_fandClaude Opus 5.5 21f30549a7 rugnux: follow a collection from the broker while it is written
rugnux <path>/<prefix>_master.h5 http://<broker>:<port>/ - when the master file is not on disk yet -
processes the broker's current collection as it is written:

- the start message comes from GET /live/start.cbor; the given path has to end in the collection's
  master file name, and the data files are read next to it;
- JFJochStartMessageReader builds the dataset from the start message (detector, geometry,
  goniometer, mask) and reads the images straight from the data files; a read of an image whose file
  the broker has not reported waits for it;
- BrokerFeed reads GET /live/events in a thread of its own and hands the reader each closed file. On
  end, the run finishes with the images the collection delivered (a read past them returns no image).
  A broken stream, another collection, or a refused follower makes every waiting read throw
  LiveCollectionLost, which ends the run like a fatal resource error;
- the token is taken from JUNGFRAUJOCH_HTTP_TOKEN (the variable the viewer uses) and sent only in
  the Authorization header; without it, and with no master file on disk, rugnux stops at once;
- rotation data and --mode mx only; plain http only (TLS is not linked into rugnux).

With the master file on disk, the URL is ignored and the run is offline as before. The streamed run
does not read the master file and is not compared with it: its result can differ from an offline run
of the same files in the last bits of the metadata.

Tests ([Live]): LiveCollection, the server-sent-event parser, the FileWriter callback, the TCP
pusher's closed-file reports with two real writers, and JFJochStartMessageReader against the
master the writer produced.

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

57 lines
2.3 KiB
C++

// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#pragma once
#include <condition_variable>
#include <mutex>
#include <stdexcept>
#include <string>
#include <vector>
#include "JFJochReader.h"
#include "HDF5ImageSource.h"
// The broker that reported a collection's data files can no longer be heard from, and an image this
// reader was asked for has not been reported: it will never be read.
class LiveCollectionLost : public std::runtime_error {
public:
explicit LiveCollectionLost(const std::string &what) : std::runtime_error(what) {}
};
// Reads a collection while it is being written: the dataset comes from its start message, as the
// broker sends it, and the images straight from the data files - /entry/data/data in each - once the
// broker reports a file closed and renamed into place (FileClosed). A read of an image whose file has
// not been reported waits for it; after Ended it returns false for an image no file holds, and after
// Lost it throws. The master file is never read.
class JFJochStartMessageReader : public JFJochReader {
public:
// master_path: where the collection's master file will be; the data files are next to it.
void Open(const StartMessage &start, const std::string &master_path);
void FileClosed(uint64_t file_number, uint64_t image_count);
void Ended();
void Lost(const std::string &why);
[[nodiscard]] uint64_t GetNumberOfImages() const override;
void Close() override;
bool ReadRawImage(int64_t image_number, JFJochReaderRawImage &image) override;
std::vector<SpotToSave> ReadSpots(int64_t image) const override { return {}; }
private:
bool LoadImage_i(std::shared_ptr<JFJochReaderDataset> &dataset, DataMessage &message,
std::vector<uint8_t> &buffer, int64_t image_number, bool update_dataset) override;
// Waits until the file of this image is reported; false when the collection ended without it.
bool WaitForImage(int64_t image_number);
HDF5ImageSource image_source_;
uint64_t number_of_images_ = 0;
uint64_t images_per_file_ = 1;
std::mutex m_;
std::condition_variable cv_;
std::vector<int64_t> file_images_; // per data file: images it holds, -1 while not reported
bool ended_ = false;
std::string lost_;
};