1#![allow(unused_imports)]
2use crate::numeric::BackendNumericPolicy;
5
6const DEVICE_WORK_QUEUE_NUMERIC: BackendNumericPolicy =
7 BackendNumericPolicy::new("device work queue");
8
9#[derive(Clone, Copy, Debug, Eq, PartialEq)]
11pub enum WorkQueueHostSync {
12 FinalOnly,
14 HostParticipates,
16}
17
18#[derive(Clone, Copy, Debug, Eq, PartialEq)]
20pub struct DeviceWorkQueueProfile {
21 pub initial_items: u64,
23 pub queue_capacity: u64,
25 pub entry_bytes: u64,
27 pub control_bytes: u64,
29 pub budget_bytes: u64,
31 pub host_sync: WorkQueueHostSync,
33}
34
35#[derive(Clone, Copy, Debug, Eq, PartialEq)]
38pub struct DeviceWorkQueueExpansionProfile {
39 pub initial_items: u64,
41 pub expansion_items: u64,
44 pub entry_bytes: u64,
46 pub control_bytes: u64,
48 pub budget_bytes: u64,
50 pub host_sync: WorkQueueHostSync,
52}
53
54#[derive(Clone, Copy, Debug, Eq, PartialEq)]
56pub struct DeviceWorkQueuePlan {
57 pub queue_bytes: u64,
59 pub control_bytes: u64,
61 pub resident_bytes: u64,
63 pub initial_occupancy_bps: u32,
65 pub final_only_host_sync: bool,
67}
68
69#[derive(Clone, Copy, Debug, Eq, PartialEq)]
71pub enum DeviceWorkQueueDrainStrategy {
72 SingleResidentDrain,
74 ChunkedResidentDrain,
77}
78
79#[derive(Clone, Copy, Debug, Eq, PartialEq)]
81pub struct DeviceWorkQueueBackpressurePlan {
82 pub queue: DeviceWorkQueuePlan,
84 pub strategy: DeviceWorkQueueDrainStrategy,
86 pub items_per_chunk: u64,
88 pub chunks: u64,
90 pub final_only_host_sync: bool,
92}
93
94#[derive(Clone, Debug, Eq, PartialEq)]
96pub enum DeviceWorkQueueError {
97 ZeroCapacity,
99 ZeroEntryBytes,
101 ZeroDrainChunk,
103 InitialItemsExceedCapacity {
105 initial_items: u64,
107 queue_capacity: u64,
109 },
110 HostParticipationRejected,
112 ByteCountOverflow {
114 field: &'static str,
116 },
117 OverBudget {
119 required_bytes: u64,
121 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
179pub 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
223pub 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
258pub 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}