Skip to main content

sim_lib_compute_wgpu/
queue.rs

1//! Bounded submission queue planning for wgpu tensor work.
2
3/// Snapshot of queued wgpu tensor submissions.
4#[derive(Clone, Debug, Default, PartialEq, Eq)]
5pub struct WgpuQueueSnapshot {
6    /// Queued node count.
7    pub nodes: usize,
8    /// Queued byte count.
9    pub bytes: u64,
10}
11
12/// Queue bounds for resident tensor submissions.
13#[derive(Clone, Debug, PartialEq, Eq)]
14pub struct WgpuQueueLimits {
15    /// Maximum queued nodes.
16    pub max_nodes: usize,
17    /// Maximum queued bytes.
18    pub max_bytes: u64,
19    /// Maximum accepted deadline tick.
20    pub deadline_tick: u64,
21}
22
23/// Deterministic bounded queue used before command submission.
24#[derive(Clone, Debug)]
25pub struct WgpuSubmissionQueue {
26    limits: WgpuQueueLimits,
27    nodes: usize,
28    bytes: u64,
29}
30
31impl WgpuSubmissionQueue {
32    /// Builds an empty bounded queue.
33    pub fn new(limits: WgpuQueueLimits) -> Self {
34        Self {
35            limits,
36            nodes: 0,
37            bytes: 0,
38        }
39    }
40
41    /// Accepts one submission when node, byte, and deadline bounds hold.
42    pub fn push(&mut self, bytes: u64, deadline_tick: u64) -> Result<(), String> {
43        if deadline_tick > self.limits.deadline_tick {
44            return Err("wgpu submission deadline expired".to_owned());
45        }
46        if self.nodes >= self.limits.max_nodes {
47            return Err("wgpu submission queue node limit reached".to_owned());
48        }
49        if self.bytes.saturating_add(bytes) > self.limits.max_bytes {
50            return Err("wgpu submission queue byte limit reached".to_owned());
51        }
52        self.nodes += 1;
53        self.bytes += bytes;
54        Ok(())
55    }
56
57    /// Clears synchronized submissions.
58    pub fn flush(&mut self) -> WgpuQueueSnapshot {
59        let snapshot = self.snapshot();
60        self.nodes = 0;
61        self.bytes = 0;
62        snapshot
63    }
64
65    /// Returns current queued pressure.
66    pub fn snapshot(&self) -> WgpuQueueSnapshot {
67        WgpuQueueSnapshot {
68            nodes: self.nodes,
69            bytes: self.bytes,
70        }
71    }
72}