Skip to main content

vyre_driver/
device_work_queue.rs

1#![allow(unused_imports)]
2//! Backend-neutral device-side work queue planning for dependent dataflow execution.
3
4use crate::numeric::BackendNumericPolicy;
5
6const DEVICE_WORK_QUEUE_NUMERIC: BackendNumericPolicy =
7    BackendNumericPolicy::new("device work queue");
8
9/// Host synchronization policy for a device device-side work queue.
10#[derive(Clone, Copy, Debug, Eq, PartialEq)]
11pub enum WorkQueueHostSync {
12    /// Host reads only final completion state after device-side draining.
13    FinalOnly,
14    /// Host participates during queue draining.
15    HostParticipates,
16}
17
18/// Work queue workload profile.
19#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct DeviceWorkQueueProfile {
21    /// Initial active work items enqueued before launch.
22    pub initial_items: u64,
23    /// Maximum resident queue capacity in work items.
24    pub queue_capacity: u64,
25    /// ABI bytes per queue entry.
26    pub entry_bytes: u64,
27    /// Bytes required for queue head/tail counters and changed flags.
28    pub control_bytes: u64,
29    /// Caller-approved device-memory budget.
30    pub budget_bytes: u64,
31    /// Host synchronization policy.
32    pub host_sync: WorkQueueHostSync,
33}
34
35/// Work queue profile where a resident queue should reserve device-side
36/// expansion headroom in addition to the initial frontier.
37#[derive(Clone, Copy, Debug, Eq, PartialEq)]
38pub struct DeviceWorkQueueExpansionProfile {
39    /// Initial active work items enqueued before launch.
40    pub initial_items: u64,
41    /// Additional device-produced work items the queue should absorb when the
42    /// explicit queue budget leaves enough room.
43    pub expansion_items: u64,
44    /// ABI bytes per queue entry.
45    pub entry_bytes: u64,
46    /// Bytes required for queue head/tail counters and changed flags.
47    pub control_bytes: u64,
48    /// Caller-approved device-memory budget for the resident queue.
49    pub budget_bytes: u64,
50    /// Host synchronization policy.
51    pub host_sync: WorkQueueHostSync,
52}
53
54/// Device-side work queue execution plan.
55#[derive(Clone, Copy, Debug, Eq, PartialEq)]
56pub struct DeviceWorkQueuePlan {
57    /// Resident queue bytes.
58    pub queue_bytes: u64,
59    /// Resident control bytes.
60    pub control_bytes: u64,
61    /// Total resident bytes.
62    pub resident_bytes: u64,
63    /// Queue occupancy in basis points before device-side expansion.
64    pub initial_occupancy_bps: u32,
65    /// Whether the plan guarantees final-state-only host synchronization.
66    pub final_only_host_sync: bool,
67}
68
69/// Device-side work queue drain strategy.
70#[derive(Clone, Copy, Debug, Eq, PartialEq)]
71pub enum DeviceWorkQueueDrainStrategy {
72    /// One resident drain window covers the whole queue.
73    SingleResidentDrain,
74    /// Queue capacity is split into multiple resident drain windows to bound
75    /// per-launch queue pressure without host participation.
76    ChunkedResidentDrain,
77}
78
79/// Device-side work queue plan with bounded resident drain windows.
80#[derive(Clone, Copy, Debug, Eq, PartialEq)]
81pub struct DeviceWorkQueueBackpressurePlan {
82    /// Base resident queue byte plan.
83    pub queue: DeviceWorkQueuePlan,
84    /// Selected resident drain strategy.
85    pub strategy: DeviceWorkQueueDrainStrategy,
86    /// Maximum queue entries drained by one device-side window.
87    pub items_per_chunk: u64,
88    /// Number of resident drain windows required to cover queue capacity.
89    pub chunks: u64,
90    /// Whether the backpressure plan preserves final-state-only host sync.
91    pub final_only_host_sync: bool,
92}
93
94/// Device work queue planning errors.
95#[derive(Clone, Debug, Eq, PartialEq)]
96pub enum DeviceWorkQueueError {
97    /// Queue capacity must be non-zero.
98    ZeroCapacity,
99    /// Entry ABI width must be explicit and non-zero.
100    ZeroEntryBytes,
101    /// Device-side drain chunk size must be non-zero.
102    ZeroDrainChunk,
103    /// Initial queue contents exceed capacity.
104    InitialItemsExceedCapacity {
105        /// Initial active items.
106        initial_items: u64,
107        /// Queue capacity.
108        queue_capacity: u64,
109    },
110    /// Host participation would reintroduce CPU orchestration.
111    HostParticipationRejected,
112    /// Byte arithmetic overflowed.
113    ByteCountOverflow {
114        /// Field being computed.
115        field: &'static str,
116    },
117    /// Queue does not fit the explicit device budget.
118    OverBudget {
119        /// Required bytes.
120        required_bytes: u64,
121        /// Budget bytes.
122        budget_bytes: u64,
123    },
124}
125
126impl std::fmt::Display for DeviceWorkQueueError {
127    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
128        match self {
129            Self::ZeroCapacity => write!(
130                f,
131                "device work queue capacity is zero. Fix: size the resident queue before launch."
132            ),
133            Self::ZeroEntryBytes => write!(
134                f,
135                "device work queue entry_bytes is zero. Fix: pass the concrete queue-entry ABI width."
136            ),
137            Self::ZeroDrainChunk => write!(
138                f,
139                "device work queue drain chunk is zero. Fix: pass a non-zero device-side drain window."
140            ),
141            Self::InitialItemsExceedCapacity {
142                initial_items,
143                queue_capacity,
144            } => write!(
145                f,
146                "device work queue initial_items={initial_items} exceeds queue_capacity={queue_capacity}. Fix: shard initial frontier items or increase explicit queue capacity."
147            ),
148            Self::HostParticipationRejected => write!(
149                f,
150                "device work queue rejected host participation. Fix: use final-only completion readback so dependent dataflow stays device-side."
151            ),
152            Self::ByteCountOverflow { field } => write!(
153                f,
154                "device work queue overflowed while computing {field}. Fix: shard the dependent dataflow workload before queue planning."
155            ),
156            Self::OverBudget {
157                required_bytes,
158                budget_bytes,
159            } => write!(
160                f,
161                "device work queue requires {required_bytes} bytes but budget allows {budget_bytes}. Fix: reduce queue capacity, shard the graph, or raise the explicit device budget."
162            ),
163        }
164    }
165}
166
167impl std::error::Error for DeviceWorkQueueError {}
168
169fn checked_add(lhs: u64, rhs: u64, field: &'static str) -> Result<u64, DeviceWorkQueueError> {
170    lhs.checked_add(rhs)
171        .ok_or(DeviceWorkQueueError::ByteCountOverflow { field })
172}
173
174fn checked_mul(lhs: u64, rhs: u64, field: &'static str) -> Result<u64, DeviceWorkQueueError> {
175    lhs.checked_mul(rhs)
176        .ok_or(DeviceWorkQueueError::ByteCountOverflow { field })
177}
178
179/// Plan a device-resident work queue for dependent dataflow execution.
180pub fn plan_device_work_queue(
181    profile: DeviceWorkQueueProfile,
182) -> Result<DeviceWorkQueuePlan, DeviceWorkQueueError> {
183    if profile.queue_capacity == 0 {
184        return Err(DeviceWorkQueueError::ZeroCapacity);
185    }
186    if profile.entry_bytes == 0 {
187        return Err(DeviceWorkQueueError::ZeroEntryBytes);
188    }
189    if profile.initial_items > profile.queue_capacity {
190        return Err(DeviceWorkQueueError::InitialItemsExceedCapacity {
191            initial_items: profile.initial_items,
192            queue_capacity: profile.queue_capacity,
193        });
194    }
195    if profile.host_sync != WorkQueueHostSync::FinalOnly {
196        return Err(DeviceWorkQueueError::HostParticipationRejected);
197    }
198
199    let queue_bytes = checked_mul(profile.queue_capacity, profile.entry_bytes, "queue bytes")?;
200    let resident_bytes = checked_add(queue_bytes, profile.control_bytes, "resident bytes")?;
201    if resident_bytes > profile.budget_bytes {
202        return Err(DeviceWorkQueueError::OverBudget {
203            required_bytes: resident_bytes,
204            budget_bytes: profile.budget_bytes,
205        });
206    }
207    let initial_occupancy_bps = DEVICE_WORK_QUEUE_NUMERIC.ratio_basis_points_u64(
208        profile.initial_items,
209        profile.queue_capacity,
210        0,
211        "device work queue initial occupancy",
212    );
213
214    Ok(DeviceWorkQueuePlan {
215        queue_bytes,
216        control_bytes: profile.control_bytes,
217        resident_bytes,
218        initial_occupancy_bps,
219        final_only_host_sync: true,
220    })
221}
222
223/// Plan a device-resident work queue that preserves initial-frontier capacity
224/// and uses remaining queue budget for device-side expansion headroom.
225pub fn plan_device_work_queue_with_expansion(
226    profile: DeviceWorkQueueExpansionProfile,
227) -> Result<DeviceWorkQueuePlan, DeviceWorkQueueError> {
228    let desired_capacity = checked_add(
229        profile.initial_items,
230        profile.expansion_items,
231        "queue expansion capacity",
232    )?;
233    if profile.entry_bytes == 0 {
234        return plan_device_work_queue(DeviceWorkQueueProfile {
235            initial_items: profile.initial_items,
236            queue_capacity: desired_capacity,
237            entry_bytes: profile.entry_bytes,
238            control_bytes: profile.control_bytes,
239            budget_bytes: profile.budget_bytes,
240            host_sync: profile.host_sync,
241        });
242    }
243    let budget_capacity =
244        profile.budget_bytes.saturating_sub(profile.control_bytes) / profile.entry_bytes;
245    let queue_capacity = desired_capacity
246        .min(budget_capacity)
247        .max(profile.initial_items);
248    plan_device_work_queue(DeviceWorkQueueProfile {
249        initial_items: profile.initial_items,
250        queue_capacity,
251        entry_bytes: profile.entry_bytes,
252        control_bytes: profile.control_bytes,
253        budget_bytes: profile.budget_bytes,
254        host_sync: profile.host_sync,
255    })
256}
257
258/// Plan a device-resident work queue plus bounded device-side drain windows.
259pub fn plan_device_work_queue_backpressure(
260    profile: DeviceWorkQueueProfile,
261    max_items_per_drain_launch: u64,
262) -> Result<DeviceWorkQueueBackpressurePlan, DeviceWorkQueueError> {
263    if max_items_per_drain_launch == 0 {
264        return Err(DeviceWorkQueueError::ZeroDrainChunk);
265    }
266    let queue = plan_device_work_queue(profile)?;
267    let chunks = div_ceil_u64(
268        profile.queue_capacity,
269        max_items_per_drain_launch,
270        "drain chunks",
271    )?;
272    let strategy = if chunks == 1 {
273        DeviceWorkQueueDrainStrategy::SingleResidentDrain
274    } else {
275        DeviceWorkQueueDrainStrategy::ChunkedResidentDrain
276    };
277    Ok(DeviceWorkQueueBackpressurePlan {
278        queue,
279        strategy,
280        items_per_chunk: max_items_per_drain_launch.min(profile.queue_capacity),
281        chunks,
282        final_only_host_sync: true,
283    })
284}
285
286fn div_ceil_u64(lhs: u64, rhs: u64, field: &'static str) -> Result<u64, DeviceWorkQueueError> {
287    DEVICE_WORK_QUEUE_NUMERIC
288        .checked_ceil_div_u64(lhs, rhs)
289        .ok_or(DeviceWorkQueueError::ByteCountOverflow { field })
290}
291
292#[cfg(test)]
293mod tests {
294    use super::*;
295
296    #[test]
297    fn device_work_queue_plans_final_only_resident_execution() {
298        let plan = plan_device_work_queue(DeviceWorkQueueProfile {
299            initial_items: 256,
300            queue_capacity: 1_024,
301            entry_bytes: 16,
302            control_bytes: 128,
303            budget_bytes: 32_768,
304            host_sync: WorkQueueHostSync::FinalOnly,
305        })
306        .expect("Fix: valid device work queue should plan");
307
308        assert_eq!(plan.queue_bytes, 16_384);
309        assert_eq!(plan.control_bytes, 128);
310        assert_eq!(plan.resident_bytes, 16_512);
311        assert_eq!(plan.initial_occupancy_bps, 2_500);
312        assert!(plan.final_only_host_sync);
313    }
314
315    #[test]
316    fn device_work_queue_expansion_uses_budgeted_resident_headroom() {
317        let plan = plan_device_work_queue_with_expansion(DeviceWorkQueueExpansionProfile {
318            initial_items: 4,
319            expansion_items: 12,
320            entry_bytes: 8,
321            control_bytes: 64,
322            budget_bytes: 256,
323            host_sync: WorkQueueHostSync::FinalOnly,
324        })
325        .expect("Fix: expansion headroom should fit inside the explicit queue budget");
326
327        assert_eq!(plan.queue_bytes, 128);
328        assert_eq!(plan.control_bytes, 64);
329        assert_eq!(plan.resident_bytes, 192);
330        assert_eq!(
331            plan.initial_occupancy_bps, 2_500,
332            "Fix: occupancy must use the expanded resident queue capacity"
333        );
334        assert!(plan.final_only_host_sync);
335    }
336
337    #[test]
338    fn device_work_queue_expansion_clamps_to_budget_without_dropping_initial_items() {
339        let plan = plan_device_work_queue_with_expansion(DeviceWorkQueueExpansionProfile {
340            initial_items: 4,
341            expansion_items: 100,
342            entry_bytes: 8,
343            control_bytes: 16,
344            budget_bytes: 96,
345            host_sync: WorkQueueHostSync::FinalOnly,
346        })
347        .expect("Fix: queue expansion should use all affordable headroom");
348
349        assert_eq!(plan.queue_bytes, 80);
350        assert_eq!(plan.resident_bytes, 96);
351        assert_eq!(
352            plan.initial_occupancy_bps, 4_000,
353            "Fix: initial occupancy should reflect budget-clamped expansion capacity"
354        );
355    }
356
357    #[test]
358    fn device_work_queue_expansion_fails_when_initial_frontier_cannot_fit() {
359        assert_eq!(
360            plan_device_work_queue_with_expansion(DeviceWorkQueueExpansionProfile {
361                initial_items: 8,
362                expansion_items: 100,
363                entry_bytes: 16,
364                control_bytes: 64,
365                budget_bytes: 128,
366                host_sync: WorkQueueHostSync::FinalOnly,
367            })
368            .expect_err("initial frontier must fail when it cannot fit the explicit budget"),
369            DeviceWorkQueueError::OverBudget {
370                required_bytes: 192,
371                budget_bytes: 128,
372            }
373        );
374    }
375
376    #[test]
377    fn device_work_queue_expansion_rejects_capacity_overflow() {
378        assert_eq!(
379            plan_device_work_queue_with_expansion(DeviceWorkQueueExpansionProfile {
380                initial_items: u64::MAX,
381                expansion_items: 1,
382                entry_bytes: 1,
383                control_bytes: 0,
384                budget_bytes: u64::MAX,
385                host_sync: WorkQueueHostSync::FinalOnly,
386            })
387            .expect_err("overflowed expansion capacity must fail before queue planning"),
388            DeviceWorkQueueError::ByteCountOverflow {
389                field: "queue expansion capacity",
390            }
391        );
392    }
393
394    #[test]
395    fn device_work_queue_rejects_host_participation() {
396        assert_eq!(
397            plan_device_work_queue(DeviceWorkQueueProfile {
398                initial_items: 1,
399                queue_capacity: 8,
400                entry_bytes: 16,
401                control_bytes: 64,
402                budget_bytes: 1_024,
403                host_sync: WorkQueueHostSync::HostParticipates,
404            })
405            .expect_err("host participation should fail"),
406            DeviceWorkQueueError::HostParticipationRejected
407        );
408    }
409
410    #[test]
411    fn device_work_queue_rejects_invalid_capacity_and_budget() {
412        assert_eq!(
413            plan_device_work_queue(DeviceWorkQueueProfile {
414                initial_items: 9,
415                queue_capacity: 8,
416                entry_bytes: 16,
417                control_bytes: 64,
418                budget_bytes: 1_024,
419                host_sync: WorkQueueHostSync::FinalOnly,
420            })
421            .expect_err("initial overflow should fail"),
422            DeviceWorkQueueError::InitialItemsExceedCapacity {
423                initial_items: 9,
424                queue_capacity: 8,
425            }
426        );
427        assert_eq!(
428            plan_device_work_queue(DeviceWorkQueueProfile {
429                initial_items: 1,
430                queue_capacity: 8,
431                entry_bytes: 16,
432                control_bytes: 64,
433                budget_bytes: 128,
434                host_sync: WorkQueueHostSync::FinalOnly,
435            })
436            .expect_err("over-budget queue should fail"),
437            DeviceWorkQueueError::OverBudget {
438                required_bytes: 192,
439                budget_bytes: 128,
440            }
441        );
442    }
443
444    #[test]
445    fn device_work_queue_occupancy_uses_widened_arithmetic_for_huge_queues() {
446        let plan = plan_device_work_queue(DeviceWorkQueueProfile {
447            initial_items: u64::MAX,
448            queue_capacity: u64::MAX,
449            entry_bytes: 1,
450            control_bytes: 0,
451            budget_bytes: u64::MAX,
452            host_sync: WorkQueueHostSync::FinalOnly,
453        })
454        .expect("Fix: max-sized byte queue should fit exactly");
455
456        assert_eq!(
457            plan.initial_occupancy_bps, 10_000,
458            "Fix: device work-queue occupancy must not use saturating u64 multiplication before division; full queues must report 10000 bps even near u64::MAX."
459        );
460    }
461
462    #[test]
463    fn device_work_queue_backpressure_chunks_large_resident_queues_without_host_participation() {
464        let plan = plan_device_work_queue_backpressure(
465            DeviceWorkQueueProfile {
466                initial_items: 4_096,
467                queue_capacity: 65_536,
468                entry_bytes: 16,
469                control_bytes: 128,
470                budget_bytes: 2 << 20,
471                host_sync: WorkQueueHostSync::FinalOnly,
472            },
473            8_192,
474        )
475        .expect("Fix: large resident work queue should plan bounded device-side drain chunks");
476
477        assert_eq!(
478            plan.strategy,
479            DeviceWorkQueueDrainStrategy::ChunkedResidentDrain
480        );
481        assert_eq!(plan.items_per_chunk, 8_192);
482        assert_eq!(plan.chunks, 8);
483        assert_eq!(plan.queue.resident_bytes, 1_048_704);
484        assert!(plan.final_only_host_sync);
485        assert!(plan.queue.final_only_host_sync);
486    }
487
488    #[test]
489    fn device_work_queue_backpressure_ceil_division_handles_max_capacity() {
490        let plan = plan_device_work_queue_backpressure(
491            DeviceWorkQueueProfile {
492                initial_items: u64::MAX,
493                queue_capacity: u64::MAX,
494                entry_bytes: 1,
495                control_bytes: 0,
496                budget_bytes: u64::MAX,
497                host_sync: WorkQueueHostSync::FinalOnly,
498            },
499            65_536,
500        )
501        .expect("Fix: ceil division for max-capacity queues must not overflow");
502
503        assert_eq!(
504            plan.strategy,
505            DeviceWorkQueueDrainStrategy::ChunkedResidentDrain
506        );
507        assert_eq!(plan.queue.queue_bytes, u64::MAX);
508        assert_eq!(plan.items_per_chunk, 65_536);
509        assert_eq!(plan.chunks, 281_474_976_710_656);
510        assert!(plan.final_only_host_sync);
511    }
512
513    #[test]
514    fn device_work_queue_backpressure_rejects_zero_drain_chunk() {
515        let err = plan_device_work_queue_backpressure(
516            DeviceWorkQueueProfile {
517                initial_items: 1,
518                queue_capacity: 8,
519                entry_bytes: 16,
520                control_bytes: 64,
521                budget_bytes: 1_024,
522                host_sync: WorkQueueHostSync::FinalOnly,
523            },
524            0,
525        )
526        .expect_err("zero drain chunk must fail loudly");
527
528        assert_eq!(err, DeviceWorkQueueError::ZeroDrainChunk);
529    }
530
531    #[test]
532    fn generated_device_work_queue_profiles_preserve_budget_and_sync_contracts() {
533        let mut state = 0xa409_3822_299f_31d0_u64;
534        for case_index in 0..2048usize {
535            let queue_capacity = 1 + next_u64(&mut state) % 262_144;
536            let entry_bytes = 1 + next_u64(&mut state) % 256;
537            let initial_items = next_u64(&mut state) % (queue_capacity + 1);
538            let control_bytes = next_u64(&mut state) % 4096;
539            let queue_bytes = queue_capacity
540                .checked_mul(entry_bytes)
541                .expect("Fix: generated queue byte count should fit");
542            let resident_bytes = queue_bytes
543                .checked_add(control_bytes)
544                .expect("Fix: generated resident byte count should fit");
545            let budget_bytes = resident_bytes + (next_u64(&mut state) % 8192);
546            let profile = DeviceWorkQueueProfile {
547                initial_items,
548                queue_capacity,
549                entry_bytes,
550                control_bytes,
551                budget_bytes,
552                host_sync: WorkQueueHostSync::FinalOnly,
553            };
554
555            let plan = plan_device_work_queue(profile)
556                .expect("Fix: generated valid queue profile must plan");
557            assert_eq!(plan.queue_bytes, queue_bytes, "case {case_index}");
558            assert_eq!(plan.control_bytes, control_bytes, "case {case_index}");
559            assert_eq!(plan.resident_bytes, resident_bytes, "case {case_index}");
560            assert!(plan.resident_bytes <= budget_bytes, "case {case_index}");
561            assert!(plan.initial_occupancy_bps <= 10_000, "case {case_index}");
562            assert!(plan.final_only_host_sync, "case {case_index}");
563
564            let drain = 1 + next_u64(&mut state) % queue_capacity;
565            let backpressure = plan_device_work_queue_backpressure(profile, drain)
566                .expect("Fix: generated valid backpressure profile must plan");
567            assert_eq!(backpressure.queue, plan, "case {case_index}");
568            assert!(
569                backpressure.items_per_chunk <= queue_capacity,
570                "case {case_index}"
571            );
572            assert!(backpressure.chunks >= 1, "case {case_index}");
573            assert!(backpressure.final_only_host_sync, "case {case_index}");
574
575            let expansion_items = next_u64(&mut state) % queue_capacity;
576            let expansion_budget = resident_bytes + (expansion_items * entry_bytes);
577            let expansion =
578                plan_device_work_queue_with_expansion(DeviceWorkQueueExpansionProfile {
579                    initial_items,
580                    expansion_items,
581                    entry_bytes,
582                    control_bytes,
583                    budget_bytes: expansion_budget,
584                    host_sync: WorkQueueHostSync::FinalOnly,
585                })
586                .expect("Fix: generated valid expansion queue profile must plan");
587            assert!(
588                expansion.resident_bytes <= expansion_budget,
589                "case {case_index}"
590            );
591            assert!(
592                expansion.queue_bytes >= initial_items * entry_bytes,
593                "case {case_index}"
594            );
595            assert!(expansion.final_only_host_sync, "case {case_index}");
596        }
597    }
598
599    fn next_u64(state: &mut u64) -> u64 {
600        *state = state
601            .wrapping_mul(6_364_136_223_846_793_005)
602            .wrapping_add(1_442_695_040_888_963_407);
603        *state
604    }
605}