载入中...
搜索中...
未找到
OnnxCompiler.h
浏览该文件的文档.
1#pragma once
2#include <algorithm>
3#include <atomic>
4#include <chrono>
5#include <condition_variable>
6#include <deque>
7#include <future>
8#include <mutex>
9#include <thread>
11
13// Session-owned CPU workers. Jobs own their source and never capture Graphics,
14// Run, buffers or callbacks. Joining the workers needs no device access.
17 std::vector<uint32_t> words;
18 double milliseconds = 0;
19};
22 using Job = std::packaged_task<CompiledProgram()>;
23 std::mutex mutex;
24 std::condition_variable ready;
25 std::deque<Job> jobs;
26 bool stopping = false;
27 size_t inFlight = 0;
28 std::vector<std::thread> workers;
29 std::atomic<size_t> running{0}, peak{0};
30 void stop() noexcept {
31 {
32 std::lock_guard lock(mutex);
33 stopping = true;
34 }
35 ready.notify_all();
36 for (auto& worker : workers)
37 if (worker.joinable()) worker.join();
38 }
39
40public:
42 explicit CompilerQueue(uint32_t count) {
43 try {
44 for (uint32_t i = 0; i < count; ++i)
45 workers.emplace_back([this] {
46 for (;;) {
47 Job job;
48 {
50 std::unique_lock lock(mutex);
51 ready.wait(lock, [this] { return stopping || !jobs.empty(); });
52 if (jobs.empty()) return;
53 job = std::move(jobs.front());
54 jobs.pop_front();
55 ++inFlight;
56 }
57 const auto n = ++running;
58 auto p = peak.load();
59 while (p < n && !peak.compare_exchange_weak(p, n)) {
60 }
62 job(); // packaged_task transports failures to the owning device thread.
63 --running;
64 {
66 std::lock_guard lock(mutex);
67 --inFlight;
68 }
69 ready.notify_all();
70 }
71 });
72 } catch (...) {
74 stop();
75 throw;
76 }
77 }
79 ~CompilerQueue() { stop(); }
80 // Called on the device thread. Run bounds each segment to 48 operations;
81 // this independent queue bound also protects accidental future callers.
83 std::future<CompiledProgram> enqueue(std::string source) {
85 Job job([source = std::move(source)] {
86 const auto start = std::chrono::steady_clock::now();
88 if (!result.ok()) throw std::runtime_error(result.error()->message() + "\nCompute source:\n" + source);
89 return CompiledProgram{
91 std::move(result.value()),
92 std::chrono::duration<double, std::milli>(std::chrono::steady_clock::now() - start).count()};
93 });
94 auto future = job.get_future();
95 {
97 std::lock_guard lock(mutex);
98 if (stopping || jobs.size() >= 48) throw std::runtime_error("ONNX compiler queue unavailable");
99 jobs.push_back(std::move(job));
100 }
101 ready.notify_one();
102 return future;
103 }
105 size_t peakWorkers() const noexcept { return peak.load(); }
107 void waitIdle() noexcept {
109 std::unique_lock lock(mutex);
110 ready.wait(lock, [this] { return jobs.empty() && inFlight == 0; });
111 }
113 void resetPeak() noexcept { peak = 0; }
114};
115} // namespace eve::tensor::onnx_detail
Duration start
std::mutex mutex
std::uint32_t count
const UnitySourceAsset & source
size_t peakWorkers() const noexcept
Peak workers.
~CompilerQueue()
Releases CompilerQueue resources.
void waitIdle() noexcept
Waits idle.
std::future< CompiledProgram > enqueue(std::string source)
Enqueue.
void resetPeak() noexcept
Resets peak.
CompilerQueue(uint32_t count)
Constructs a CompilerQueue.
EVENGINE_API_WORLD Result< std::vector< uint32_t > > compileComputeSpirv(const std::string &source)
Compile GLSL compute source to owning SPIR-V without accessing Graphics or a GPU.
Definition Gpgpu.cpp:294
std::future< vkb::Instance > future
Definition Graphics.cpp:182