载入中...
搜索中...
未找到
OnnxGpgpu.cpp
浏览该文件的文档.
1#include <algorithm>
2#include <chrono>
3#include <cstdlib>
4#include <iostream>
5#include <map>
6#include <unordered_map>
9#include "gpgpu/Gpgpu.h"
10#include "gpgpu/GpuBuffer.h"
11#include "gpgpu/Sequence.h"
12#include "graphics/Graphics.h"
13#include "tensor/OnnxCompiler.h"
14#include "tensor/OnnxCompute.h"
15namespace eve::tensor {
16namespace {
17using Clock = std::chrono::steady_clock;
18struct Timed {
19 double& milliseconds;
20 Clock::time_point start = Clock::now();
21 ~Timed() { milliseconds += std::chrono::duration<double, std::milli>(Clock::now() - start).count(); }
22};
23struct Run;
24struct DeviceBytes final : OnnxDeviceStorage {
25 std::weak_ptr<Run> run;
26 std::unique_ptr<gpgpu::GpuBuffer> buffer;
27 size_t size = 0, capacity = 0;
28 uint64_t epoch = 0;
29 ~DeviceBytes() override;
30 Result<std::vector<uint8_t>> readback() const override;
31};
32struct ShaderEntry {
33 std::future<onnx_detail::CompiledProgram> compilation;
34 std::unique_ptr<gpgpu::ComputeShader> shader;
35};
36struct PendingDispatch {
37 ShaderEntry* program;
38 std::vector<std::shared_ptr<DeviceBytes>> bindings;
39 uint32_t groups;
40};
41struct Run : std::enable_shared_from_this<Run> {
42 gpgpu::Gpgpu& gp;
43 std::unique_ptr<gpgpu::Sequence> sequence;
44 std::unordered_map<std::string, ShaderEntry> shaders;
45 onnx_detail::CompilerQueue compiler;
46 std::vector<PendingDispatch> dispatches;
48 std::map<std::shared_ptr<const std::vector<uint8_t>>, std::shared_ptr<DeviceBytes>,
49 std::owner_less<std::shared_ptr<const std::vector<uint8_t>>>>
51 std::vector<std::shared_ptr<const std::vector<uint8_t>>> uncommittedUploads;
52 bool active = true;
53 uint64_t epoch = 1;
54 std::vector<std::shared_ptr<DeviceBytes>> pending;
55 std::vector<std::weak_ptr<DeviceBytes>> resources;
56 OnnxTransferStats stats;
58 size_t compiledShaders = 0;
59 Clock::time_point started = Clock::now();
60 size_t queued = 0, cachedBytes = 0, allocatedBytes = 0;
61 std::multimap<size_t, std::unique_ptr<gpgpu::GpuBuffer>> freeBuffers;
62 explicit Run(gpgpu::Gpgpu& g, uint32_t workers) : gp(g), sequence(g.newSequence()), compiler(workers) {
63 sequence->begin();
64 }
65 ~Run() {
66 // Destroy an unsubmitted command sequence before its referenced resources.
67 sequence.reset();
68 for (auto& [key, entry] : shaders)
69 if (entry.shader) entry.shader->clearBindings();
70 for (auto& weak : resources)
71 if (auto p = weak.lock()) p->buffer.reset();
72 }
73 void beginCall() noexcept {
74 stats = {};
75 compiler.resetPeak();
79 started = Clock::now();
80 active = true;
81 ++epoch;
82 }
83 void endCall() noexcept {
84 compiler.waitIdle();
85 // An abandoned graph must not record or submit its deferred commands.
86 // Unmaterialized jobs contain only CPU data; discard their future (including any
87 // compilation exception) because the caller already received the graph failure.
88 dispatches.clear();
89 std::erase_if(shaders, [](const auto& item) { return !item.second.shader; });
90 if (const auto* profile = std::getenv("EVE_ONNX_PROFILE"); profile && std::string(profile) == "1")
91 std::cerr << "ONNX_PROFILE shader_ms=" << shaderMs << " allocation_ms=" << allocationMs
92 << " upload_record_ms=" << uploadRecordMs << " dispatch_record_ms=" << dispatchRecordMs
93 << " submit_wait_ms=" << waitMs << " readback_host_ms=" << readbackMs
94 << " compile_task_ms=" << compileTaskMs << " compile_wait_ms=" << compileWaitMs
95 << " pipeline_ms=" << pipelineMs << " compiler_peak_workers=" << compiler.peakWorkers()
96 << " compiled_shaders=" << compiledShaders << " resident_scope_ms="
97 << std::chrono::duration<double, std::milli>(Clock::now() - started).count() << '\n';
98 sequence.reset();
99 for (auto& [key, entry] : shaders)
100 if (entry.shader) entry.shader->clearBindings();
101 // Abandoned recorded uploads have not reached the device and cannot enter a persistent cache.
102 for (const auto& key : uncommittedUploads) {
103 if (auto it = persistentUploads.find(key); it != persistentUploads.end()) {
104 cachedBytes -= key->size();
105 persistentUploads.erase(it);
106 }
107 }
108 uncommittedUploads.clear();
109 pending.clear();
110 uploads.clear();
111 queued = 0;
112 active = false;
113 std::erase_if(resources, [](const auto& p) { return p.expired(); });
114 }
115 std::shared_ptr<DeviceBytes> allocate(size_t bytes) {
116 Timed timing{allocationMs};
117 if (!bytes || bytes > INT32_MAX - 3) throw std::runtime_error("Invalid resident buffer extent");
118 auto p = std::make_shared<DeviceBytes>();
119 p->run = shared_from_this();
120 p->size = bytes;
121 p->epoch = epoch;
122 const auto capacity = std::max(size_t(4), (bytes + 3) & ~size_t(3));
123 auto reusable = freeBuffers.lower_bound(capacity);
124 if (reusable != freeBuffers.end() && reusable->first <= std::max(size_t(4096), capacity * 2)) {
125 p->capacity = reusable->first;
126 p->buffer = std::move(reusable->second);
127 freeBuffers.erase(reusable);
128 ++stats.bufferReuses;
129 } else {
130 if (allocatedBytes + capacity > 512u * 1024u * 1024u) {
131 for (const auto& [n, b] : freeBuffers) allocatedBytes -= n;
132 freeBuffers.clear();
133 }
134 if (allocatedBytes + capacity > 512u * 1024u * 1024u)
135 throw std::runtime_error("ONNX GPU resident memory limit exceeded");
136 p->buffer.reset(gp.newBuffer(static_cast<int>(capacity)));
137 p->capacity = capacity;
139 ++stats.bufferAllocations;
140 }
141 resources.push_back(p);
142 return p;
143 }
144 void recordDispatches() {
145 for (auto& command : dispatches) {
146 auto& entry = *command.program;
147 if (!entry.shader) {
148 Timed shaderTiming{shaderMs};
149 onnx_detail::CompiledProgram program;
150 {
151 Timed waitTiming{compileWaitMs};
152 program = entry.compilation.get();
153 }
154 compileTaskMs += program.milliseconds;
155 {
156 Timed pipelineTiming{pipelineMs};
157 auto result = gpgpu::createComputeShader(program.words);
158 if (!result.ok()) throw std::runtime_error(result.error()->message());
159 entry.shader = std::move(result.value());
160 }
161 }
162 Timed timing{dispatchRecordMs};
163 auto& shader = *entry.shader;
164 for (size_t i = 0; i < command.bindings.size(); ++i)
165 shader.bindBuffer(static_cast<int>(i), command.bindings[i]->buffer.get());
166 sequence->recordDispatch(&shader, static_cast<int>(command.groups));
167 shader.clearBindings();
168 }
169 stats.compilerPeakWorkers = compiler.peakWorkers();
170 dispatches.clear();
171 }
172 void flush() {
173 if (!queued) return;
174 recordDispatches();
175 Timed timing{waitMs};
176 // Submit asynchronously through the existing sequence contract; CPU visibility is
177 // requested only at a graph boundary. Retain every referenced allocation until wait.
178 if (sequence->submitAsync() != gpgpu::SequenceStatus::Submitted ||
180 throw std::runtime_error("ONNX GPU submission failed");
181 uncommittedUploads.clear();
182 ++stats.submissions;
183 pending.clear();
184 queued = 0;
185 sequence->begin();
186 }
187 std::shared_ptr<DeviceBytes> input(const OnnxBuffer& value) {
188 if (value.device) {
189 auto d = std::dynamic_pointer_cast<DeviceBytes>(value.device);
190 if (!d || d->run.lock().get() != this || !d->buffer || d->epoch != epoch || !active)
191 throw std::runtime_error("Foreign or expired GPU storage");
192 return d;
193 }
194 if (!value.host || value.host->size() != value.size) throw std::runtime_error("Invalid host input storage");
195 auto& cache = value.persistent ? persistentUploads : uploads;
196 if (auto it = cache.find(value.host); it != cache.end()) {
197 it->second->epoch = epoch;
198 return it->second;
199 }
200 if (value.persistent && cachedBytes + value.size > 128u * 1024u * 1024u) {
201 persistentUploads.clear();
202 cachedBytes = 0;
203 }
204 auto d = allocate(std::max(size_t(1), value.size));
205 Timed timing{uploadRecordMs};
206 std::vector<uint8_t> padded(std::max(size_t(4), (value.size + 3) & ~size_t(3)), 0);
207 std::copy(value.host->begin(), value.host->end(), padded.begin());
208 sequence->recordUpload(d->buffer.get(), padded.data(), padded.size());
209 ++stats.uploads;
210 stats.uploadedBytes += value.size;
211 ++queued;
212 cache.emplace(value.host, d);
213 if (value.persistent) {
214 cachedBytes += value.size;
215 uncommittedUploads.push_back(value.host);
216 }
217 pending.push_back(d);
218 return d;
219 }
220};
221DeviceBytes::~DeviceBytes() {
222 if (buffer)
223 if (auto owner = run.lock()) owner->freeBuffers.emplace(capacity, std::move(buffer));
224}
225Result<std::vector<uint8_t>> DeviceBytes::readback() const {
226 try {
227 auto owner = run.lock();
228 if (!owner || !buffer || !owner->active || epoch != owner->epoch)
229 throw std::runtime_error("ONNX GPU storage expired");
230 // Dispatches must precede the download in the command stream. Compilation
231 // waits are accounted separately from host readback and GPU submission.
232 owner->recordDispatches();
233 const auto start = Clock::now();
234 const auto priorWait = owner->waitMs;
235 std::unique_ptr<gpgpu::GpuBuffer> staging(owner->gp.newBuffer(static_cast<int>(size), "staging"));
236 owner->sequence->recordDownload(buffer.get(), staging.get(), size);
237 ++owner->queued;
238 owner->flush();
239 std::vector<uint8_t> out(size);
240 staging->downloadBytes(out.data(), size);
241 owner->readbackMs +=
242 std::chrono::duration<double, std::milli>(Clock::now() - start).count() - (owner->waitMs - priorWait);
243 ++owner->stats.downloads;
244 owner->stats.downloadedBytes += size;
245 return Result<std::vector<uint8_t>>::success(std::move(out));
246 } catch (const std::exception& e) {
247 return Result<std::vector<uint8_t>>::failure(Diagnostic::error(DiagnosticCode::Failed, e.what()));
248 }
249}
250struct SessionState {
251 std::shared_ptr<Run> run;
252 bool retired = false;
253 uint32_t compilerWorkers = 4;
254};
255class EngineCompute final : public OnnxCompute {
256 std::shared_ptr<SessionState> state;
257 eve::Subscription retirement;
258 OnnxTransferStats last;
259 bool active = false;
260
261protected:
262 void beginRun() noexcept override {
263 if (state->run) state->run->beginCall();
264 last = {};
265 active = true;
266 }
267 void endRun() noexcept override {
268 if (state->run) {
269 last = state->run->stats;
270 state->run->endCall();
271 }
272 active = false;
273 }
274
275public:
276 EngineCompute(std::shared_ptr<SessionState> s, eve::Subscription token)
277 : state(std::move(s)), retirement(std::move(token)) {}
278 OnnxTransferStats transferStats() const override { return active && state->run ? state->run->stats : last; }
279 Result<OnnxBuffer> enqueue(const std::string& source, std::span<const OnnxBuffer> inputs, size_t bytes,
280 uint32_t work) override {
281 try {
282 auto& run = state->run;
283 if (state->retired) throw std::runtime_error("ONNX GPU session device has been retired");
284 if (!active) throw std::runtime_error("Resident enqueue requires an active runGpu scope");
285 if (inputs.size() >= 8 || !bytes || !work || (uint64_t(work) + 63) / 64 > 65535)
286 throw std::runtime_error("Invalid GPU kernel extent");
287 if (!run) {
288 auto* gp = gpgpu::Gpgpu::create();
289 if (!gp || !gp->isAvailable()) throw std::runtime_error("ONNX GPU device unavailable");
290 run = std::make_shared<Run>(*gp, state->compilerWorkers);
291 }
292 if (!run->sequence) {
293 run->sequence.reset(run->gp.newSequence());
294 run->sequence->begin();
295 }
296 // Bound pending resources and the existing shader descriptor pool (64 sets).
297 if (run->queued >= 48) run->flush();
298 auto it = run->shaders.find(source);
299 if (it == run->shaders.end()) {
300 if (run->shaders.size() >= 1024) {
301 run->flush();
302 run->shaders.clear();
303 }
304 ++run->stats.shaderCompilations;
305 ++run->compiledShaders;
306 it = run->shaders.emplace(source, ShaderEntry{run->compiler.enqueue(source), {}}).first;
307 } else
308 ++run->stats.shaderCacheHits;
309 PendingDispatch command{&it->second, {}, (work + 63) / 64};
310 for (const auto& input : inputs) {
311 auto d = run->input(input);
312 command.bindings.push_back(d);
313 run->pending.push_back(d);
314 }
315 auto out = run->allocate(bytes);
316 command.bindings.push_back(out);
317 run->dispatches.push_back(std::move(command));
318 run->pending.push_back(out);
319 ++run->queued;
320 return Result<OnnxBuffer>::success({{}, out, bytes});
321 } catch (const std::exception& e) {
323 Diagnostic::error(DiagnosticCode::Failed, e.what(), {}, {}, "tensor.onnx.gpu"));
324 }
325 }
326 Result<std::vector<uint8_t>> dispatch(const OnnxKernel& k) override {
327 const bool temporary = !active;
328 if (temporary) beginRun();
329 struct Scope {
330 EngineCompute& owner;
331 bool temporary;
332 ~Scope() {
333 if (temporary) owner.endRun();
334 }
335 } scope{*this, temporary};
336 std::vector<OnnxBuffer> inputs;
337 for (auto s : k.inputs) {
338 auto p = std::make_shared<const std::vector<uint8_t>>(s.begin(), s.end());
339 inputs.push_back({p, {}, p->size()});
340 }
341 auto r = enqueue(k.source, inputs, k.outputBytes, k.workItems);
342 if (!r.ok()) return Result<std::vector<uint8_t>>::failure(r.status());
343 return r.value().device->readback();
344 }
345};
346} // namespace
348 if (compilerWorkers < 1 || compilerWorkers > 8)
350 Diagnostic::error(DiagnosticCode::InvalidArgument, "Compiler worker count must be 1..8"));
351#ifdef EVENGINE_WEBGPU
353 Diagnostic::error(DiagnosticCode::Unsupported, "ONNX GLSL kernels require the Vulkan backend"));
354#else
355 try {
356 auto* gp = gpgpu::Gpgpu::create();
357 if (gp && gp->isAvailable()) {
358 auto state = std::make_shared<SessionState>();
359 state->compilerWorkers = compilerWorkers;
360 std::weak_ptr<SessionState> weak = state;
361 auto token = graphics::Graphics::create()->onResourcesRetiring([weak]() noexcept {
362 if (auto s = weak.lock()) {
363 s->retired = true;
364 s->run.reset();
365 }
366 });
367 if (!token.ok()) return Result<std::unique_ptr<OnnxCompute>>::failure(token.status());
369 std::make_unique<EngineCompute>(state, std::move(token.value())));
370 }
371 } catch (const std::exception& e) {
373 }
375 DiagnosticCode::Unsupported, "Initialize the engine Vulkan graphics device before ONNX GPU execution"));
376#endif
377}
378} // namespace eve::tensor
double value
Duration start
bool & active
const std::string & s
std::vector< QuestEvent > pending
glm::vec4 p[6]
EvpackChunkInput input
Definition Evpack.cpp:170
tensor::Graph g
Definition GpuGraph.cpp:7
bool retired
std::uint32_t capacity
std::uint32_t key
glm::vec3 n
Definition Grass.cpp:63
int inputs
Definition GridGraph.cpp:23
double r
std::int32_t first
JobScope scope
std::uint64_t bytes
MeleePoint3 b
Definition MeleeHit.cpp:41
uint32_t groups
Definition OnnxGpgpu.cpp:39
gpgpu::Gpgpu & gp
Definition OnnxGpgpu.cpp:42
size_t compiledShaders
Definition OnnxGpgpu.cpp:58
std::vector< std::weak_ptr< DeviceBytes > > resources
Definition OnnxGpgpu.cpp:55
double shaderMs
Definition OnnxGpgpu.cpp:57
std::future< onnx_detail::CompiledProgram > compilation
Definition OnnxGpgpu.cpp:33
double allocationMs
Definition OnnxGpgpu.cpp:57
std::map< std::shared_ptr< const std::vector< uint8_t > >, std::shared_ptr< DeviceBytes >, std::owner_less< std::shared_ptr< const std::vector< uint8_t > > > > persistentUploads
Definition OnnxGpgpu.cpp:50
std::vector< std::shared_ptr< DeviceBytes > > bindings
Definition OnnxGpgpu.cpp:38
std::weak_ptr< Run > run
Definition OnnxGpgpu.cpp:25
double dispatchRecordMs
Definition OnnxGpgpu.cpp:57
double pipelineMs
Definition OnnxGpgpu.cpp:47
onnx_detail::CompilerQueue compiler
Definition OnnxGpgpu.cpp:45
size_t allocatedBytes
Definition OnnxGpgpu.cpp:60
uint64_t epoch
Definition OnnxGpgpu.cpp:28
double waitMs
Definition OnnxGpgpu.cpp:57
OnnxTransferStats stats
Definition OnnxGpgpu.cpp:56
std::map< std::shared_ptr< const std::vector< uint8_t > >, std::shared_ptr< DeviceBytes >, std::owner_less< std::shared_ptr< const std::vector< uint8_t > > > > uploads
Definition OnnxGpgpu.cpp:50
double readbackMs
Definition OnnxGpgpu.cpp:57
std::unique_ptr< gpgpu::GpuBuffer > buffer
Definition OnnxGpgpu.cpp:26
double compileTaskMs
Definition OnnxGpgpu.cpp:47
std::unordered_map< std::string, ShaderEntry > shaders
Definition OnnxGpgpu.cpp:44
size_t cachedBytes
Definition OnnxGpgpu.cpp:60
size_t queued
Definition OnnxGpgpu.cpp:60
std::multimap< size_t, std::unique_ptr< gpgpu::GpuBuffer > > freeBuffers
Definition OnnxGpgpu.cpp:61
std::unique_ptr< gpgpu::Sequence > sequence
Definition OnnxGpgpu.cpp:43
double compileWaitMs
Definition OnnxGpgpu.cpp:47
double & milliseconds
Definition OnnxGpgpu.cpp:19
uint32_t compilerWorkers
std::vector< PendingDispatch > dispatches
Definition OnnxGpgpu.cpp:46
double uploadRecordMs
Definition OnnxGpgpu.cpp:57
std::vector< std::shared_ptr< const std::vector< uint8_t > > > uncommittedUploads
Definition OnnxGpgpu.cpp:51
float d
Shader * shader
uint64_t token
std::shared_ptr< const ExpressionProgram > program
float size
Definition TreeMesh.cpp:156
const UnitySourceAsset & source
static Diagnostic error(DiagnosticCode code, std::string message, std::string path={}, DiagnosticDetails details={}, std::string source={})
Construct an error diagnostic with the standard error severity.
Definition Diagnostic.h:125
Move-only operation result carrying either a value or Status.
Definition Result.h:155
static Result success(T value)
Construct a successful result owning value.
Definition Result.h:164
static Result failure(Status status)
Construct a failed result from a structured status.
Definition Result.h:175
Move-only RAII token that owns one observer registration.
EVENGINE_API_WORLD Result< std::unique_ptr< ComputeShader > > createComputeShader(const std::vector< uint32_t > &words)
Create an owning compute pipeline from valid SPIR-V produced by compileComputeSpirv.
Definition Gpgpu.cpp:307
EVENGINE_API_DOMAINS Result< std::unique_ptr< OnnxCompute > > createOnnxGpuCompute(uint32_t compilerWorkers=4)
Create a reusable GPU session for the active engine Gpgpu Vulkan device, or an error.
bool started
Definition Graphics.cpp:183