From aa3fa6b9c76e6983ed1afa440059168dbef21f09 Mon Sep 17 00:00:00 2001 From: Filip Leonarski Date: Sun, 27 Sep 2026 10:32:58 +0200 Subject: [PATCH] Reader: stream the frame data into the page cache ahead of the image loops A cold run on a spinning disk waited on the disk twice over. The CBF header scan read 256 kB from every frame on eight threads that each strode through their own share of the sweep, so they drifted apart and the scan became a seek storm (34 s for 2400 frames here); and after it, the pre-scan and the first-pass indexing touch a few hundred frames and leave the disk idle until the first image loop reads everything at seek-bound rates. - ReadAhead (reader/): once the dataset is open, rugnux starts eight threads that read the data files - HDF5 data files (legacy, VDS or the integrated master) or the per-frame CBF/marCCD/SMV files - in 4 MB pieces taken strictly in order, into a throwaway buffer. One stream reads this disk at 125 MB/s, eight in-order streams at 190 MB/s, 32 at 157 MB/s. It never gets more than a quarter of MemAvailable (GlobalMemoryStatusEx on Windows, 4 GiB where there is no figure) ahead of what ReadRawImage has handed out, so a dataset bigger than the cache does not evict its own start, and it stops with the reader. Plain ifstream reads: portable, no POSIX calls. - Header scans (CBF, marCCD, SMV) hand the files out in order from an atomic counter (sweep::ForEachInOrder) instead of striding: 18 s -> 12 s for 2400 cold CBF headers. The CBF header is first read with a 16 kB probe and again with the old 256 kB one only when the separator is not in it, so the parsed header is exactly what it was: 12 s -> 6 s. Output unchanged: p.hkl, p.mtz and p_unmerged.mtz md5-identical to the rc173 baseline on 6toc (CBF, 2400 frames, 6.0 GB) and 9q41 (HDF5 VDS, 900 frames, 5.1 GB), and on 6z9g (HDF5, 12.8 GB) to the unmodified branch; myob (p.hkl p.mtz p_P1.mtz p_unmerged.mtz) md5-identical to the reference. Measured cold (files evicted with POSIX_FADV_DONTNEED before every run), same code without this commit vs with it, on a shared box (load 20-70, other agents reading the same disk, so single runs scatter by +-20 s): 6toc wall 61.7/62.7 -> 49.3/49.8 s (clean pairs); all data resident after 62/51/50 -> 45/41/42 s 9q41 wall 67.3 -> 57.8 s (clean pair); resident after 58/46/43 -> 48/37/35 s 6z9g resident after 81 -> 69 s The first image loop can look slower with this in CBF runs: the old 256 kB header probes pulled ~70% of the data in as kernel readahead, so the old loop started warm - after a 34 s header scan instead of 10 s. Warm (myob, NVMe, cached): 19.35/20.00 s without, 19.76-20.16 s with; the read-ahead then only copies 9.3 GB out of the page cache, 0.44 s wall and 3.4 CPU-s measured standalone. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01D1G8gJVAy6gp1K5Dz3NE5C --- reader/CMakeLists.txt | 2 + reader/HDF5ImageLocator.cpp | 18 ++++++++ reader/HDF5ImageLocator.h | 4 ++ reader/HDF5ImageSource.h | 3 ++ reader/JFJochCBFReader.cpp | 22 +++------- reader/JFJochCBFReader.h | 2 + reader/JFJochHDF5Reader.cpp | 7 +++ reader/JFJochHDF5Reader.h | 2 + reader/JFJochMarCCDReader.cpp | 22 +++------- reader/JFJochMarCCDReader.h | 2 + reader/JFJochReader.cpp | 7 +++ reader/JFJochReader.h | 16 +++++++ reader/JFJochSMVReader.cpp | 22 +++------- reader/JFJochSMVReader.h | 2 + reader/MiniCBF.cpp | 12 +++++- reader/ReadAhead.cpp | 81 +++++++++++++++++++++++++++++++++++ reader/ReadAhead.h | 58 +++++++++++++++++++++++++ reader/SweepLayout.cpp | 16 +++++++ reader/SweepLayout.h | 8 ++++ rugnux/rugnux_cli.cpp | 4 ++ tests/JFJochReaderTest.cpp | 15 +++++++ 21 files changed, 275 insertions(+), 50 deletions(-) create mode 100644 reader/ReadAhead.cpp create mode 100644 reader/ReadAhead.h diff --git a/reader/CMakeLists.txt b/reader/CMakeLists.txt index 06f304776..d2488fa27 100644 --- a/reader/CMakeLists.txt +++ b/reader/CMakeLists.txt @@ -1,5 +1,7 @@ ADD_LIBRARY(JFJochReader STATIC JFJochReader.cpp JFJochReader.h + ReadAhead.cpp + ReadAhead.h JFJochHDF5Reader.cpp JFJochHDF5Reader.h SweepLayout.cpp diff --git a/reader/HDF5ImageLocator.cpp b/reader/HDF5ImageLocator.cpp index 05d6f908e..eed666c85 100644 --- a/reader/HDF5ImageLocator.cpp +++ b/reader/HDF5ImageLocator.cpp @@ -4,6 +4,8 @@ #include "HDF5ImageLocator.h" #include "../common/JFJochException.h" +#include + namespace { // Coalesce consecutive single-image mappings into one contiguous range when the source and // virtual images stay contiguous in the same file/dataset. @@ -101,6 +103,22 @@ HDF5ImageLocator::Location HDF5ImageLocator::Resolve(int64_t global_image) const return {layout_.master_file, static_cast(global_image), layout_.master_filename}; } +std::vector HDF5ImageLocator::DataFiles() const { + std::vector ret; + if (layout_.format == FileWriterFormat::NXmxLegacy) { + for (const auto &f: layout_.legacy_files) + ret.push_back(f.path); + } else if (layout_.format == FileWriterFormat::NXmxVDS + && layout_.data_layout == HDF5DataSetLayout::VIRTUAL) { + for (const auto &mapping: layout_.vds_mappings) + if (CoversFirstChannel(mapping) + && std::find(ret.begin(), ret.end(), mapping.filename) == ret.end()) + ret.push_back(mapping.filename); + } else if (!layout_.master_filename.empty()) + ret.push_back(layout_.master_filename); + return ret; +} + std::vector HDF5ImageLocator::GetSourceMapping(uint64_t first_image, std::optional image_count, uint64_t total_images, diff --git a/reader/HDF5ImageLocator.h b/reader/HDF5ImageLocator.h index bc5bf6579..c08377e7a 100644 --- a/reader/HDF5ImageLocator.h +++ b/reader/HDF5ImageLocator.h @@ -64,6 +64,10 @@ public: // by the layout. Does not bounds-check against the total image count - the caller does that. Location Resolve(int64_t global_image) const; + // The files that hold the pixels, in image order: the data files of a legacy or VDS dataset, the + // master itself when the images are in it. + std::vector DataFiles() const; + // Source mapping for re-writing a derived file (e.g. _process.h5) so it links back to the // original pixel sources rather than to a master. total_images is supplied by the caller. // stride is the step between consecutive images of the derived file in the SOURCE: image i of the diff --git a/reader/HDF5ImageSource.h b/reader/HDF5ImageSource.h index b3bc77dff..fe53c17f1 100644 --- a/reader/HDF5ImageSource.h +++ b/reader/HDF5ImageSource.h @@ -70,6 +70,9 @@ public: // file that holds a legacy/VDS image's per-image metadata. HDF5ImageLocator::Location Resolve(int64_t global) const; + // The files holding the pixels, in image order. + std::vector DataFiles() const { return locator_.DataFiles(); } + // Read the pixels at a resolved location into a CompressedImage backed by `buffer`. Templated on // the allocator so a caller can hand over a buffer that does not zero what it is about to // overwrite (RawByteBuffer). diff --git a/reader/JFJochCBFReader.cpp b/reader/JFJochCBFReader.cpp index 2f69f93ff..117bea505 100644 --- a/reader/JFJochCBFReader.cpp +++ b/reader/JFJochCBFReader.cpp @@ -8,10 +8,8 @@ #include #include #include -#include #include #include -#include #include #include "../common/JFJochException.h" @@ -299,20 +297,11 @@ void JFJochCBFReader::ReadFiles(const std::string &path) { // sweep of several thousand frames this is seconds of startup before a single image is read, and // the files are independent. std::vector frames(files_.size()); - { - const size_t nthreads = std::min(std::max(1u, std::thread::hardware_concurrency()), 8); - std::vector> futures; - for (size_t t = 0; t < nthreads; t++) - futures.push_back(std::async(std::launch::async, [&, t] { - for (size_t i = t; i < files_.size(); i += nthreads) { - const auto h = minicbf::ReadHeader(files_[i]); - frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m, - h.beam_x_px, h.beam_y_px, h.wavelength_A}; - } - })); - for (auto &f : futures) - f.get(); - } + sweep::ForEachInOrder(files_.size(), [&](size_t i) { + const auto h = minicbf::ReadHeader(files_[i]); + frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m, + h.beam_x_px, h.beam_y_px, h.wavelength_A}; + }); // The sweep the headers describe, which is not always the files laid out end to end: a deposited // series can be missing frames, and those are gaps in the rotation rather than images to close up. @@ -406,6 +395,7 @@ bool JFJochCBFReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &i if (!HasImage(image_number)) return false; image.image = DecodeInto(image_number, image.image_buffer, image.read_buffer); + NoteImageRead(); return true; } diff --git a/reader/JFJochCBFReader.h b/reader/JFJochCBFReader.h index 1ee73ab89..0511fad47 100644 --- a/reader/JFJochCBFReader.h +++ b/reader/JFJochCBFReader.h @@ -37,6 +37,8 @@ class JFJochCBFReader : public JFJochReader { template CompressedImage DecodeInto(int64_t image_number, Buffer &buffer, std::vector &scratch) const; + std::vector DataFiles() const override { return files_; } + public: ~JFJochCBFReader() override = default; diff --git a/reader/JFJochHDF5Reader.cpp b/reader/JFJochHDF5Reader.cpp index 4009a8cdf..bffbc0296 100644 --- a/reader/JFJochHDF5Reader.cpp +++ b/reader/JFJochHDF5Reader.cpp @@ -72,14 +72,21 @@ bool JFJochHDF5Reader::ReadRawImage(int64_t image_number, JFJochReaderRawImage & chunk = image_source_.PrepareDirectRead(loc); if (!chunk) { ret.image = image_source_.ReadImageAt(ret.image_buffer, loc); + NoteImageRead(); return true; } } ret.image = HDF5ImageSource::ReadDirect(ret.image_buffer, *chunk); + NoteImageRead(); return true; } +std::vector JFJochHDF5Reader::DataFiles() const { + std::unique_lock ul(hdf5_mutex); + return image_source_.DataFiles(); +} + bool JFJochHDF5Reader::LoadImage_i(std::shared_ptr &dataset, DataMessage &message, std::vector &buffer, diff --git a/reader/JFJochHDF5Reader.h b/reader/JFJochHDF5Reader.h index 175fedd89..6f2650464 100644 --- a/reader/JFJochHDF5Reader.h +++ b/reader/JFJochHDF5Reader.h @@ -40,6 +40,8 @@ class JFJochHDF5Reader : public JFJochReader { HDF5ImageLocator::Location GetImageLocation(int64_t image_number) const; + std::vector DataFiles() const override; + public: ~JFJochHDF5Reader() override = default; diff --git a/reader/JFJochMarCCDReader.cpp b/reader/JFJochMarCCDReader.cpp index 49acfe551..552b67aa1 100644 --- a/reader/JFJochMarCCDReader.cpp +++ b/reader/JFJochMarCCDReader.cpp @@ -4,8 +4,6 @@ #include "JFJochMarCCDReader.h" #include -#include -#include #include "../common/JFJochException.h" #include "../common/JFJochMath.h" @@ -116,20 +114,11 @@ void JFJochMarCCDReader::ReadFiles(const std::string &path) { // Reading one costs a 4 kB read, so on a sweep of several thousand frames this is worth // spreading over the cores, as the CBF path does for the same reason. std::vector frames(files_.size()); - { - const size_t nthreads = std::min(std::max(1u, std::thread::hardware_concurrency()), 8); - std::vector> futures; - for (size_t t = 0; t < nthreads; t++) - futures.push_back(std::async(std::launch::async, [&, t] { - for (size_t i = t; i < files_.size(); i += nthreads) { - const auto h = marccd::ReadHeader(files_[i]); - frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m, - h.beam_x_px, h.beam_y_px, h.wavelength_A}; - } - })); - for (auto &f : futures) - f.get(); - } + sweep::ForEachInOrder(files_.size(), [&](size_t i) { + const auto h = marccd::ReadHeader(files_[i]); + frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m, + h.beam_x_px, h.beam_y_px, h.wavelength_A}; + }); // The sweep the headers describe, which is not always the files laid out end to end: a deposited // series can be missing frames, and those are gaps in the rotation rather than images to close @@ -215,6 +204,7 @@ bool JFJochMarCCDReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage if (!HasImage(image_number)) return false; image.image = DecodeInto(image_number, image.image_buffer, image.read_buffer); + NoteImageRead(); return true; } diff --git a/reader/JFJochMarCCDReader.h b/reader/JFJochMarCCDReader.h index 886009162..74aa7be0c 100644 --- a/reader/JFJochMarCCDReader.h +++ b/reader/JFJochMarCCDReader.h @@ -38,6 +38,8 @@ class JFJochMarCCDReader : public JFJochReader { template CompressedImage DecodeInto(int64_t image_number, Buffer &buffer, std::vector &scratch) const; + std::vector DataFiles() const override { return files_; } + public: ~JFJochMarCCDReader() override = default; diff --git a/reader/JFJochReader.cpp b/reader/JFJochReader.cpp index e8912ce0a..f7a17a073 100644 --- a/reader/JFJochReader.cpp +++ b/reader/JFJochReader.cpp @@ -64,6 +64,13 @@ std::shared_ptr JFJochReader::GetRawImage(int64_t image_nu return ret; } +void JFJochReader::StartReadAhead() { + read_ahead_.reset(); + const auto files = DataFiles(); + if (!files.empty() && GetNumberOfImages() > 0) + read_ahead_ = std::make_unique(files, GetNumberOfImages()); +} + void JFJochReader::SetStartMessage(const std::shared_ptr &val) { std::unique_lock ul(m); dataset = val; diff --git a/reader/JFJochReader.h b/reader/JFJochReader.h index 31a5b59d8..d59075ecf 100644 --- a/reader/JFJochReader.h +++ b/reader/JFJochReader.h @@ -6,6 +6,7 @@ #include #include +#include #include #include "../common/JFJochMessages.h" @@ -15,6 +16,7 @@ #include "JFJochReaderDataset.h" #include "JFJochReaderImage.h" #include "JFJochReaderSpots.h" +#include "ReadAhead.h" class JFJochReader { mutable std::mutex m; @@ -31,8 +33,17 @@ class JFJochReader { int64_t n_image, int64_t image_jump, JFJochReaderImage &image); + + std::unique_ptr read_ahead_; protected: void SetStartMessage(const std::shared_ptr &val); + + // The files that hold the frame data, in the order the images sit in them: what StartReadAhead() + // streams. Empty for a reader with no files of its own. + virtual std::vector DataFiles() const { return {}; } + // Every ReadRawImage() that read an image reports it here, which is what moves the read-ahead on. + void NoteImageRead() { if (read_ahead_) read_ahead_->ImageRead(); } + DiffractionExperiment default_experiment; public: virtual ~JFJochReader() = default; @@ -44,6 +55,11 @@ public: virtual void Close() = 0; + // Start streaming the frame data into the page cache ahead of the image reads (see ReadAhead). + // For a batch run over the whole dataset; call it once the dataset is open, before any image is + // read. It changes only when the bytes arrive, never what is read. + void StartReadAhead(); + std::shared_ptr LoadImage(int64_t image_number, int64_t summation_factor = 1); void UpdateGeomMetadata(const DiffractionExperiment& experiment); diff --git a/reader/JFJochSMVReader.cpp b/reader/JFJochSMVReader.cpp index 2100538cd..548a746e6 100644 --- a/reader/JFJochSMVReader.cpp +++ b/reader/JFJochSMVReader.cpp @@ -4,8 +4,6 @@ #include "JFJochSMVReader.h" #include -#include -#include #include "../common/JFJochException.h" #include "../common/JFJochMath.h" @@ -116,20 +114,11 @@ void JFJochSMVReader::ReadFiles(const std::string &path) { // Reading one costs a 4 kB read, so on a sweep of several thousand frames this is worth // spreading over the cores, as the CBF path does for the same reason. std::vector frames(files_.size()); - { - const size_t nthreads = std::min(std::max(1u, std::thread::hardware_concurrency()), 8); - std::vector> futures; - for (size_t t = 0; t < nthreads; t++) - futures.push_back(std::async(std::launch::async, [&, t] { - for (size_t i = t; i < files_.size(); i += nthreads) { - const auto h = smv::ReadHeader(files_[i]); - frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m, - h.beam_x_px, h.beam_y_px, h.wavelength_A}; - } - })); - for (auto &f : futures) - f.get(); - } + sweep::ForEachInOrder(files_.size(), [&](size_t i) { + const auto h = smv::ReadHeader(files_[i]); + frames[i] = {files_[i], h.start_angle_deg, h.angle_increment_deg, h.distance_m, + h.beam_x_px, h.beam_y_px, h.wavelength_A}; + }); // The sweep the headers describe, which is not always the files laid out end to end: a deposited // series can be missing frames, and those are gaps in the rotation rather than images to close @@ -215,6 +204,7 @@ bool JFJochSMVReader::ReadRawImage(int64_t image_number, JFJochReaderRawImage &i if (!HasImage(image_number)) return false; image.image = DecodeInto(image_number, image.image_buffer, image.read_buffer); + NoteImageRead(); return true; } diff --git a/reader/JFJochSMVReader.h b/reader/JFJochSMVReader.h index ec67a8187..815a08e60 100644 --- a/reader/JFJochSMVReader.h +++ b/reader/JFJochSMVReader.h @@ -38,6 +38,8 @@ class JFJochSMVReader : public JFJochReader { template CompressedImage DecodeInto(int64_t image_number, Buffer &buffer, std::vector &scratch) const; + std::vector DataFiles() const override { return files_; } + public: ~JFJochSMVReader() override = default; diff --git a/reader/MiniCBF.cpp b/reader/MiniCBF.cpp index 822557121..a7bdb24fa 100644 --- a/reader/MiniCBF.cpp +++ b/reader/MiniCBF.cpp @@ -437,14 +437,22 @@ void Slurp(const std::string &path, size_t max_bytes, std::vector &buf) } // Enough to reach the separator on any header seen in the wild (the longest measured is ~6.3 kB). +// The first read is kept small because a sweep's startup reads every frame's header: on a spinning +// disk 256 kB from each of 2400 files took twice as long as 16 kB (12 s against 6 s). A header +// that does not end within it is read again with the large probe, so the result never depends on it. +constexpr size_t HEADER_FIRST_PROBE_BYTES = 16 * 1024; constexpr size_t HEADER_PROBE_BYTES = 256 * 1024; } // namespace Header ReadHeader(const std::string &path) { std::vector buf; - Slurp(path, HEADER_PROBE_BYTES, buf); - const auto start = FindBinarySection(buf.data(), buf.size()); + Slurp(path, HEADER_FIRST_PROBE_BYTES, buf); + auto start = FindBinarySection(buf.data(), buf.size()); + if (!start.has_value()) { + Slurp(path, HEADER_PROBE_BYTES, buf); + start = FindBinarySection(buf.data(), buf.size()); + } if (!start.has_value()) throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "No CBF binary section in " + path); diff --git a/reader/ReadAhead.cpp b/reader/ReadAhead.cpp new file mode 100644 index 000000000..3f87bd654 --- /dev/null +++ b/reader/ReadAhead.cpp @@ -0,0 +1,81 @@ +// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute +// SPDX-License-Identifier: GPL-3.0-only + +#include "ReadAhead.h" + +#include +#include +#include +#include + +#ifdef _WIN32 +#include +#endif + +namespace { + +constexpr uint64_t PIECE_BYTES = 4 << 20; +constexpr size_t THREADS = 8; + +// Memory the system can hand out without swapping. Where there is no such figure (macOS), 4 GiB. +uint64_t AvailableMemory() { +#ifdef _WIN32 + MEMORYSTATUSEX status{}; + status.dwLength = sizeof(status); + if (GlobalMemoryStatusEx(&status)) + return status.ullAvailPhys; +#else + std::ifstream meminfo("/proc/meminfo"); + std::string key; + uint64_t kb = 0; + while (meminfo >> key >> kb) { + if (key == "MemAvailable:") + return kb * 1024; + meminfo.ignore(std::numeric_limits::max(), '\n'); + } +#endif + return 4ULL << 30; +} + +} // namespace + +ReadAhead::ReadAhead(std::vector files, uint64_t n_images) : files_(std::move(files)) { + uint64_t total = 0; + for (size_t i = 0; i < files_.size(); i++) { + std::error_code ec; + const uint64_t size = std::filesystem::file_size(files_[i], ec); // a gap in the sweep is "" + if (ec) + continue; + for (uint64_t offset = 0; offset < size; offset += PIECE_BYTES) + pieces_.push_back({i, offset, total + offset}); + total += size; + } + bytes_per_image_ = static_cast(total) / static_cast(n_images); + window_ = static_cast(AvailableMemory() / 4); + + for (size_t t = 0; t < THREADS; t++) + threads_.emplace_back(&ReadAhead::Run, this); +} + +ReadAhead::~ReadAhead() { + stop_ = true; + for (auto &t: threads_) + t.join(); +} + +void ReadAhead::Run() { + std::vector buffer(PIECE_BYTES); + for (size_t i = next_piece_++; i < pieces_.size(); i = next_piece_++) { + const auto &piece = pieces_[i]; + while (static_cast(piece.position) > window_ + static_cast(images_read_) * bytes_per_image_) { + if (stop_) + return; + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + } + if (stop_) + return; + std::ifstream in(files_[piece.file], std::ios::binary); + in.seekg(static_cast(piece.offset)); + in.read(buffer.data(), static_cast(buffer.size())); + } +} diff --git a/reader/ReadAhead.h b/reader/ReadAhead.h new file mode 100644 index 000000000..07cae9fac --- /dev/null +++ b/reader/ReadAhead.h @@ -0,0 +1,58 @@ +// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute +// SPDX-License-Identifier: GPL-3.0-only + +#pragma once + +#include +#include +#include +#include +#include + +// Streams the files that hold a dataset's frames into the page cache, front to back, ahead of the +// image loops that read them. +// +// On a spinning disk a cold run otherwise waits on the disk: the header scan, the pre-scan and the +// first-pass indexing touch a few hundred frames and leave the disk idle for tens of seconds, and +// then the first image loop reads everything at whatever rate its scattered requests get. Streaming +// from the moment the dataset is open keeps the disk busy through all of that, at its sequential +// rate, so the loops find their frames already in memory. +// +// The files are read in 4 MB pieces by a few threads that take the pieces strictly in order: on the +// disk measured here one thread streams at 125 MB/s, eight such threads at 190 MB/s (and 32 at +// 157 MB/s), so eight keep the stream ahead of an image loop that is itself reading from disk. +// +// It changes only WHEN bytes arrive, never what is read: the loops still read every image themselves, +// and a frame the read-ahead has not reached yet is read cold, as it always was. +// +// It never runs more than a window - a quarter of the memory available when it starts - ahead of +// what the loops have read. A dataset larger than that is not streamed in only to push out its own +// beginning before the first loop gets there; and a page the read-ahead brought in is then read a +// second time by the loop, which is what keeps it in the page cache for the loops after. +class ReadAhead { + struct Piece { + size_t file; + uint64_t offset; // in the file + uint64_t position; // in the whole stream + }; + + std::vector files_; + std::vector pieces_; + double window_ = 0; + double bytes_per_image_ = 0; + std::atomic next_piece_{0}; + std::atomic images_read_{0}; + std::atomic stop_{false}; + std::vector threads_; + + void Run(); +public: + // files: in the order the images sit in them; n_images: how many images they hold between them. + ReadAhead(std::vector files, uint64_t n_images); + ~ReadAhead(); + ReadAhead(const ReadAhead &) = delete; + ReadAhead &operator=(const ReadAhead &) = delete; + + // An image loop has read one image; the window moves on by one image's worth of bytes. + void ImageRead() { images_read_++; } +}; diff --git a/reader/SweepLayout.cpp b/reader/SweepLayout.cpp index 6feae1122..d43609df9 100644 --- a/reader/SweepLayout.cpp +++ b/reader/SweepLayout.cpp @@ -4,10 +4,13 @@ #include "SweepLayout.h" #include +#include #include #include #include +#include #include +#include #include "../common/JFJochException.h" #include "../common/Logger.h" @@ -281,4 +284,17 @@ Layout Place(const std::vector &frames, const std::string &logger_name) { return out; } +void ForEachInOrder(size_t n, const std::function &fn) { + const size_t nthreads = std::min(std::max(1u, std::thread::hardware_concurrency()), 8); + std::atomic next{0}; + std::vector> futures; + for (size_t t = 0; t < nthreads; t++) + futures.push_back(std::async(std::launch::async, [&] { + for (size_t i = next++; i < n; i = next++) + fn(i); + })); + for (auto &f : futures) + f.get(); +} + } // namespace sweep diff --git a/reader/SweepLayout.h b/reader/SweepLayout.h index 31767f1c6..ab9880169 100644 --- a/reader/SweepLayout.h +++ b/reader/SweepLayout.h @@ -3,6 +3,8 @@ #pragma once +#include +#include #include #include @@ -57,4 +59,10 @@ struct Layout { // logger_name is the reader's own logger, so the gap report names the format the user passed. Layout Place(const std::vector &frames, const std::string &logger_name); +// Call fn(i) for every i in [0, n) on up to 8 threads, handing the indices out in order. For reading +// every frame's header: the threads then read neighbouring files at any moment, so a spinning disk +// sees one sweep through the directory rather than eight positions drifting apart (2400 CBF headers +// from a cold spinning disk: 18 s with each thread striding through its own share, 12 s in order). +void ForEachInOrder(size_t n, const std::function &fn); + } // namespace sweep diff --git a/rugnux/rugnux_cli.cpp b/rugnux/rugnux_cli.cpp index 8eb64c648..80f21fbe5 100644 --- a/rugnux/rugnux_cli.cpp +++ b/rugnux/rugnux_cli.cpp @@ -1456,6 +1456,10 @@ static int RunRugnux(int argc, char **argv) { exit(EXIT_FAILURE); } JFJochReader &reader = *reader_ptr; + // Stream the frames in from disk from now on, so a cold or slow disk works through the setup and + // the pre-scan instead of sitting idle until the first image loop. --mode scale reads no images. + if (mode != RugnuxMode::Scale) + reader.StartReadAhead(); const auto dataset = reader.GetDataset(); if (!dataset) { diff --git a/tests/JFJochReaderTest.cpp b/tests/JFJochReaderTest.cpp index 5403c9764..ba8d01780 100644 --- a/tests/JFJochReaderTest.cpp +++ b/tests/JFJochReaderTest.cpp @@ -4128,6 +4128,21 @@ TEST_CASE("JFJochCBFReader_incomplete_header_is_refused_not_misread", "[HDF5][Fu } } +TEST_CASE("MiniCBF_header_longer_than_the_first_probe", "[HDF5][Full]") { + // The header is read with a small first probe, and again with the large one when it does not end + // inside it; what stands past the first 16 kB must come out exactly as before. + std::string filler; + for (int i = 0; i < 1000; i++) + filler += "# Comment line padding the header out past the first probe\n"; + WriteRawMiniCBF("cbflong_0001.cbf", "# Pixel_size 172e-6 m x 172e-6 m\n" + filler + + "# Wavelength 0.96864 A\n# Detector_distance 0.33161 m\n", 24, 16); + const auto h = minicbf::ReadHeader("cbflong_0001.cbf"); + CHECK(h.pixel_x_m == Catch::Approx(172e-6)); + CHECK(h.wavelength_A == Catch::Approx(0.96864)); + CHECK(h.distance_m == Catch::Approx(0.33161)); + remove("cbflong_0001.cbf"); +} + // marCCD: a TIFF whose instrument header sits in the gap between the TIFF header and the pixels. // The fixtures below are written byte for byte rather than through libtiff, because the layout IS // what the reader has to cope with - the private tag that states where the instrument header