Build Packages / Create release (push) Successful in 16s
Build Packages / build:rugnux:aarch64 (cross) (push) Successful in 8m27s
Build Packages / build:rugnux-tgz (x86_64) (push) Successful in 9m15s
Build Packages / build:viewer-tgz:cpu (push) Successful in 10m11s
Build Packages / build:viewer-tgz:cuda (push) Successful in 12m6s
Build Packages / build:rpm (rocky8_nocuda) (push) Successful in 15m44s
Build Packages / build:rpm (rocky9_nocuda) (push) Successful in 16m1s
Build Packages / build:windows:nocuda (push) Successful in 17m29s
Build Packages / build:windows:cuda (push) Successful in 19m58s
Build Packages / HDF5 consumer tests (DIALS, XDS) (push) Successful in 24m7s
Build Packages / build:rpm (ubuntu2404_nocuda) (push) Successful in 19m8s
Build Packages / build:rugnux:windows (push) Successful in 10m58s
Build Packages / build:rpm (ubuntu2204_nocuda) (push) Successful in 20m46s
Build Packages / Generate python client (push) Successful in 53s
Build Packages / build:rpm (rocky8_sls9) (push) Successful in 20m13s
Build Packages / Build documentation (push) Successful in 1m36s
Build Packages / build:rpm (rocky9_sls9) (push) Successful in 19m57s
Build Packages / build:rpm (rocky8) (push) Successful in 18m7s
Build Packages / build:rpm (rocky9) (push) Successful in 18m54s
Build Packages / build:rpm (ubuntu2204) (push) Successful in 19m32s
Build Packages / build:rpm (ubuntu2404) (push) Successful in 17m30s
Build Packages / Unit tests (push) Successful in 1h39m2s
* Fixed `jfjoch_broker` cancelling every data collection with a CUDA "out of memory" error after long operation: GPU memory no longer leaks with each collection. * Rugnux scales a rotation sweep until the per-frame scales settle instead of for a fixed three rounds, and says so when they did not - merged intensities, and the space group, resolution cut and frame rejection read off them, change accordingly; `--scaling-iterations` is now the cap on that loop (default 100). * Rugnux places every frame of a marCCD, SMV or miniCBF series at the spindle angle its own header states, so a series with missing frames, or with angles written modulo 360, is no longer read at the wrong geometry or refused. * Every rotation run writes two diagnostic files beside its reflections: `<prefix>_detector.jpg`, the detector projection with the pixel mask and the detected beam-stop shadow drawn on it, and `<prefix>_plot.txt`, one row per image. Reviewed-on: #82 Co-authored-by: Filip Leonarski <filip.leonarski@psi.ch>
146 lines
4.0 KiB
C++
146 lines
4.0 KiB
C++
// SPDX-FileCopyrightText: 2024 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
|
|
// SPDX-License-Identifier: GPL-3.0-only
|
|
|
|
#pragma once
|
|
|
|
#include <chrono>
|
|
#include <queue>
|
|
#include <mutex>
|
|
#include <condition_variable>
|
|
#include <set>
|
|
|
|
template <class T>
|
|
class ThreadSafeFIFO {
|
|
std::queue<T> queue;
|
|
std::condition_variable c_empty, c_full;
|
|
mutable std::mutex m;
|
|
const size_t max_size;
|
|
size_t max_utilization;
|
|
size_t utilization;
|
|
bool stopped = false;
|
|
public:
|
|
explicit ThreadSafeFIFO(size_t in_max_size = UINT32_MAX) : max_size(in_max_size), max_utilization(0), utilization(0) {}
|
|
|
|
// Release every waiter and make all further blocking operations return at once: puts are
|
|
// dropped, gets answer with a default-constructed element. Used when the owner of the queue is
|
|
// torn down - at that point nobody is going to drain it any more, so a producer blocked on a
|
|
// full queue would never return and the thread join in the destructor would deadlock.
|
|
void Stop() {
|
|
std::unique_lock ul(m);
|
|
stopped = true;
|
|
c_empty.notify_all();
|
|
c_full.notify_all();
|
|
}
|
|
|
|
void Clear() {
|
|
std::unique_lock ul(m);
|
|
queue = {};
|
|
utilization = 0;
|
|
max_utilization = 0;
|
|
// A producer blocked on the full queue has to be told the room it was waiting for is there;
|
|
// nothing else would wake it, as the next Get finds the queue empty and notifies no one.
|
|
c_full.notify_all();
|
|
}
|
|
|
|
bool Put(T val) {
|
|
std::unique_lock ul(m);
|
|
if (queue.size() < max_size) {
|
|
queue.push(val);
|
|
c_empty.notify_one();
|
|
utilization++;
|
|
if (utilization > max_utilization)
|
|
max_utilization = utilization;
|
|
return true;
|
|
} else
|
|
return false;
|
|
};
|
|
|
|
void PutBlocking(T val) {
|
|
std::unique_lock ul(m);
|
|
c_full.wait(ul, [&]{return stopped || (queue.size() < max_size);});
|
|
if (stopped)
|
|
return;
|
|
queue.push(val);
|
|
utilization++;
|
|
if (utilization > max_utilization)
|
|
max_utilization = utilization;
|
|
c_empty.notify_one();
|
|
};
|
|
|
|
bool PutTimeout(T val, std::chrono::milliseconds timeout) {
|
|
std::unique_lock ul(m);
|
|
if (!c_full.wait_for(ul, timeout, [&]{ return stopped || (queue.size() < max_size); }))
|
|
return false;
|
|
if (stopped)
|
|
return false;
|
|
queue.push(val);
|
|
utilization++;
|
|
if (utilization > max_utilization)
|
|
max_utilization = utilization;
|
|
c_empty.notify_one();
|
|
return true;
|
|
}
|
|
|
|
int Get(T &val) {
|
|
std::unique_lock ul(m);
|
|
|
|
if (queue.empty())
|
|
return 0;
|
|
else {
|
|
val = queue.front();
|
|
queue.pop();
|
|
c_full.notify_one();
|
|
utilization--;
|
|
return 1;
|
|
}
|
|
}
|
|
|
|
T GetBlocking() {
|
|
std::unique_lock ul(m);
|
|
c_empty.wait(ul, [&]{return stopped || !queue.empty();});
|
|
if (queue.empty())
|
|
return T{};
|
|
T tmp = queue.front();
|
|
queue.pop();
|
|
c_full.notify_one();
|
|
utilization--;
|
|
return tmp;
|
|
};
|
|
|
|
int GetTimeout(T &val, std::chrono::microseconds timeout) {
|
|
std::unique_lock ul(m);
|
|
if (queue.empty())
|
|
c_empty.wait_for(ul, timeout, [&]{return stopped || !queue.empty();});
|
|
|
|
if (queue.empty())
|
|
return 0;
|
|
else {
|
|
val = queue.front();
|
|
queue.pop();
|
|
c_full.notify_one();
|
|
utilization--;
|
|
return 1;
|
|
}
|
|
}
|
|
|
|
[[nodiscard]] size_t Size() const {
|
|
std::unique_lock ul(m);
|
|
return queue.size();
|
|
}
|
|
|
|
void ClearMaxUtilization() {
|
|
std::unique_lock ul(m);
|
|
max_utilization = utilization;
|
|
}
|
|
|
|
[[nodiscard]] size_t GetMaxUtilization() const {
|
|
std::unique_lock ul(m);
|
|
return max_utilization;
|
|
}
|
|
|
|
[[nodiscard]] size_t GetCurrentUtilization() const {
|
|
std::unique_lock ul(m);
|
|
return utilization;
|
|
}
|
|
};
|