6#include <unordered_map>
17using Clock = std::chrono::steady_clock;
20 Clock::time_point
start = Clock::now();
21 ~Timed() {
milliseconds += std::chrono::duration<double, std::milli>(Clock::now() - start).count(); }
24struct DeviceBytes final : OnnxDeviceStorage {
25 std::weak_ptr<Run>
run;
26 std::unique_ptr<gpgpu::GpuBuffer>
buffer;
29 ~DeviceBytes()
override;
30 Result<std::vector<uint8_t>> readback()
const override;
34 std::unique_ptr<gpgpu::ComputeShader>
shader;
36struct PendingDispatch {
38 std::vector<std::shared_ptr<DeviceBytes>>
bindings;
41struct Run : std::enable_shared_from_this<Run> {
44 std::unordered_map<std::string, ShaderEntry>
shaders;
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>>>>
54 std::vector<std::shared_ptr<DeviceBytes>>
pending;
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) {
71 if (auto
p = weak.lock())
p->
buffer.reset();
73 void beginCall() noexcept {
83 void endCall() noexcept {
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()
97 << std::chrono::duration<double, std::milli>(Clock::now() - started).count() <<
'\n';
113 std::erase_if(resources, [](
const auto&
p) {
return p.expired(); });
115 std::shared_ptr<DeviceBytes> allocate(
size_t bytes) {
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();
122 const auto capacity = std::max(
size_t(4), (
bytes + 3) & ~
size_t(3));
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);
128 ++
stats.bufferReuses;
130 if (allocatedBytes +
capacity > 512u * 1024u * 1024u) {
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)));
139 ++
stats.bufferAllocations;
144 void recordDispatches() {
146 auto& entry = *command.program;
149 onnx_detail::CompiledProgram
program;
152 program = entry.compilation.get();
158 if (!result.ok())
throw std::runtime_error(result.error()->message());
159 entry.shader = std::move(result.value());
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));
180 throw std::runtime_error(
"ONNX GPU submission failed");
187 std::shared_ptr<DeviceBytes>
input(
const OnnxBuffer&
value) {
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");
194 if (!
value.host ||
value.host->size() !=
value.size)
throw std::runtime_error(
"Invalid host input storage");
196 if (
auto it = cache.find(
value.host); it != cache.end()) {
197 it->second->epoch =
epoch;
200 if (
value.persistent && cachedBytes +
value.size > 128u * 1024u * 1024u) {
204 auto d = allocate(std::max(
size_t(1),
value.size));
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());
212 cache.emplace(
value.host,
d);
213 if (
value.persistent) {
221DeviceBytes::~DeviceBytes() {
223 if (
auto owner =
run.lock()) owner->freeBuffers.emplace(
capacity, std::move(
buffer));
225Result<std::vector<uint8_t>> DeviceBytes::readback()
const {
227 auto owner =
run.lock();
228 if (!owner || !
buffer || !owner->active ||
epoch != owner->epoch)
229 throw std::runtime_error(
"ONNX GPU storage expired");
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);
239 std::vector<uint8_t> out(
size);
240 staging->downloadBytes(out.data(),
size);
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) {
251 std::shared_ptr<Run>
run;
255class EngineCompute final :
public OnnxCompute {
256 std::shared_ptr<SessionState>
state;
258 OnnxTransferStats last;
262 void beginRun()
noexcept override {
267 void endRun() noexcept
override {
269 last =
state->run->stats;
270 state->run->endCall();
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 {
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");
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);
292 if (!
run->sequence) {
293 run->sequence.reset(
run->gp.newSequence());
294 run->sequence->begin();
297 if (
run->queued >= 48)
run->flush();
299 if (it ==
run->shaders.end()) {
300 if (
run->shaders.size() >= 1024) {
302 run->shaders.clear();
304 ++
run->stats.shaderCompilations;
305 ++
run->compiledShaders;
308 ++
run->stats.shaderCacheHits;
309 PendingDispatch command{&it->second, {}, (work + 63) / 64};
312 command.bindings.push_back(
d);
313 run->pending.push_back(
d);
316 command.bindings.push_back(out);
317 run->dispatches.push_back(std::move(command));
318 run->pending.push_back(out);
321 }
catch (
const std::exception& e) {
326 Result<std::vector<uint8_t>> dispatch(
const OnnxKernel& k)
override {
327 const bool temporary = !
active;
328 if (temporary) beginRun();
330 EngineCompute& owner;
333 if (temporary) owner.endRun();
335 }
scope{*
this, temporary};
336 std::vector<OnnxBuffer>
inputs;
338 auto p = std::make_shared<const std::vector<uint8_t>>(
s.begin(),
s.end());
339 inputs.push_back({
p, {},
p->size()});
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();
348 if (compilerWorkers < 1 || compilerWorkers > 8)
351#ifdef EVENGINE_WEBGPU
356 auto*
gp = gpgpu::Gpgpu::create();
357 if (
gp &&
gp->isAvailable()) {
358 auto state = std::make_shared<SessionState>();
360 std::weak_ptr<SessionState> weak =
state;
361 auto token = graphics::Graphics::create()->onResourcesRetiring([weak]()
noexcept {
362 if (
auto s = weak.lock()) {
369 std::make_unique<EngineCompute>(
state, std::move(
token.value())));
371 }
catch (
const std::exception& e) {
std::vector< QuestEvent > pending
std::vector< std::weak_ptr< DeviceBytes > > resources
std::future< onnx_detail::CompiledProgram > compilation
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
std::vector< std::shared_ptr< DeviceBytes > > bindings
onnx_detail::CompilerQueue compiler
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
std::unique_ptr< gpgpu::GpuBuffer > buffer
std::unordered_map< std::string, ShaderEntry > shaders
std::multimap< size_t, std::unique_ptr< gpgpu::GpuBuffer > > freeBuffers
std::unique_ptr< gpgpu::Sequence > sequence
std::vector< PendingDispatch > dispatches
std::vector< std::shared_ptr< const std::vector< uint8_t > > > uncommittedUploads
std::shared_ptr< const ExpressionProgram > program
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.
Move-only operation result carrying either a value or Status.
static Result success(T value)
Construct a successful result owning value.
static Result failure(Status status)
Construct a failed result from a structured status.
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.
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.