Files
Jungfraujoch/image_analysis/indexing/IndexerThreadPool.h
T

64 lines
1.9 KiB
C++

// SPDX-FileCopyrightText: 2025 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#ifndef JFJOCH_INDEXERTHREADPOOL_H
#define JFJOCH_INDEXERTHREADPOOL_H
#include <thread>
#include <mutex>
#include <condition_variable>
#include <queue>
#include <functional>
#include <future>
#include <vector>
#include <optional>
#include <memory>
#include <latch>
#include "../common/JFJochMessages.h"
#include "../common/DiffractionSpot.h"
#include "../common/DiffractionExperiment.h"
#include "../common/NUMAHWPolicy.h"
#include "Indexer.h"
class IndexerThread {
struct TaskInput {
const DiffractionExperiment &experiment;
const std::vector<Coord> &recip;
};
bool stop = false;
enum class TaskState {STARTING, IDLE, READY, COMPLETED, ERROR} state = TaskState::STARTING;
std::mutex m;
std::condition_variable c_running;
std::condition_variable c_start;
std::condition_variable c_done;
std::unique_ptr<IndexerResult> result = nullptr;
std::unique_ptr<TaskInput> task_input = nullptr;
std::thread worker_thread;
void Worker(const IndexingSettings& settings, int threadid);
public:
IndexerThread(const IndexingSettings& settings, int threadid);
~IndexerThread();
std::unique_ptr<IndexerResult> Run(const DiffractionExperiment &experiment, const std::vector<Coord> &recip);
void Finalize();
};
class IndexerThreadPool {
std::mutex m;
std::condition_variable c;
std::vector<uint8_t> worker_busy;
size_t worker_free_count;
std::vector<std::unique_ptr<IndexerThread>> tasks;
const int64_t viable_cell_min_spots;
const bool blocking;
int GetFreeWorker();
public:
IndexerThreadPool(const IndexingSettings& settings);
IndexerResult Run(const DiffractionExperiment& experiment, const std::vector<Coord>& recip);
};
#endif //JFJOCH_INDEXERTHREADPOOL_H