1use alloc::boxed::Box;
2use alloc::rc::Rc;
3use alloc::vec::{IntoIter, Vec};
4use core::sync::atomic::{AtomicU64, Ordering};
5
6use virtio_accel_transport::{
7 ChainId, DeviceChain, DeviceQueue, DriverQueue, MalformedChain, NotificationHint,
8 NotificationRecheck, PublishError, PublishErrorKind, PublishedChain, QueueConfigError,
9 QueueControl, QueueEpoch, QueueError, QueuePort, QueueSize, QueueState, ReclaimedChain,
10 UsedChain, UsedLength,
11};
12
13use crate::{DriverChain, SplitDeviceChain};
14
15#[derive(Clone, Copy, Debug, PartialEq, Eq)]
17pub enum SplitQueueInitError {
18 DescriptorLimit,
20}
21
22#[derive(Clone, Copy, Debug, PartialEq, Eq)]
24pub enum SplitQueueError {
25 AllocationFailed,
27 Busy,
29 ChainLimitExceedsQueue,
31 MalformedDriverChain(MalformedChain),
33 InvalidTestState,
35 CorruptState,
37}
38
39#[derive(Clone, Copy, Debug, PartialEq, Eq)]
41pub struct RingCounters {
42 pub available: u16,
44 pub consumed_available: u16,
46 pub used: u16,
48 pub consumed_used: u16,
50}
51
52#[derive(Debug)]
57pub struct SplitQueue {
58 state: QueueState,
59 max_chain_descriptors: u16,
60 current_epoch: Rc<AtomicU64>,
61 storage: Option<RingStorage>,
62 available_notifications: bool,
63 used_notifications: bool,
64}
65
66impl SplitQueue {
67 pub fn new(
69 max_size: QueueSize,
70 max_chain_descriptors: u16,
71 ) -> Result<Self, SplitQueueInitError> {
72 if max_chain_descriptors < 2 || max_chain_descriptors > max_size.get() {
73 return Err(SplitQueueInitError::DescriptorLimit);
74 }
75 let epoch = QueueEpoch::INITIAL;
76 Ok(Self {
77 state: QueueState::unconfigured(max_size, epoch),
78 max_chain_descriptors,
79 current_epoch: Rc::new(AtomicU64::new(epoch.get())),
80 storage: None,
81 available_notifications: true,
82 used_notifications: true,
83 })
84 }
85
86 pub const fn max_chain_descriptors(&self) -> u16 {
88 self.max_chain_descriptors
89 }
90
91 pub fn inject_available(
96 &mut self,
97 chain: DriverChain,
98 ) -> Result<PublishedChain, PublishError<DriverChain, SplitQueueError>> {
99 self.publish_inner(chain, true)
100 }
101
102 pub fn set_empty_ring_index(&mut self, index: u16) -> Result<(), QueueError<SplitQueueError>> {
104 let storage = self.storage_mut()?;
105 if storage.active_chains != 0 {
106 return Err(QueueError::Transport(SplitQueueError::Busy));
107 }
108 storage.available_index = index;
109 storage.consumed_available = index;
110 storage.used_index = index;
111 storage.consumed_used = index;
112 storage.available_count = 0;
113 storage.used_count = 0;
114 Ok(())
115 }
116
117 pub fn set_next_descriptor(&mut self, index: u16) -> Result<(), QueueError<SplitQueueError>> {
119 let storage = self.storage_mut()?;
120 if index >= storage.size.get() {
121 return Err(QueueError::Transport(SplitQueueError::InvalidTestState));
122 }
123 storage.next_descriptor = index;
124 Ok(())
125 }
126
127 pub fn ring_counters(&self) -> Option<RingCounters> {
129 self.storage.as_ref().map(RingStorage::counters)
130 }
131
132 pub fn free_descriptors(&self) -> Option<u16> {
134 self.storage
135 .as_ref()
136 .map(|storage| storage.free_descriptors)
137 }
138
139 fn publish_inner(
140 &mut self,
141 mut chain: DriverChain,
142 inject_malformed: bool,
143 ) -> Result<PublishedChain, PublishError<DriverChain, SplitQueueError>> {
144 if !self.state.ready() {
145 return Err(PublishError::new(chain, PublishErrorKind::NotReady));
146 }
147 if !inject_malformed {
148 if let Err(error) = chain.validation() {
149 return Err(PublishError::new(
150 chain,
151 PublishErrorKind::Transport(SplitQueueError::MalformedDriverChain(error)),
152 ));
153 }
154 if chain.descriptor_count() > self.max_chain_descriptors {
155 return Err(PublishError::new(
156 chain,
157 PublishErrorKind::InsufficientDescriptors,
158 ));
159 }
160 }
161
162 let epoch = self.state.epoch();
163 let notifications = self.available_notifications;
164 let storage = match self.storage.as_mut() {
165 Some(storage) => storage,
166 None => {
167 return Err(PublishError::new(
168 chain,
169 PublishErrorKind::Transport(SplitQueueError::CorruptState),
170 ));
171 }
172 };
173 if storage.available_count == storage.size.get() {
174 return Err(PublishError::new(chain, PublishErrorKind::QueueFull));
175 }
176 if chain.descriptor_count() > storage.free_descriptors {
177 return Err(PublishError::new(
178 chain,
179 PublishErrorKind::InsufficientDescriptors,
180 ));
181 }
182
183 let head = match storage.allocate_descriptors(&mut chain) {
184 Ok(head) => head,
185 Err(error) => {
186 return Err(PublishError::new(chain, PublishErrorKind::Transport(error)));
187 }
188 };
189 let id = ChainId::new(epoch, u64::from(head));
190 if storage.records[usize::from(head)].is_some() {
191 storage.release_descriptors(&chain);
192 return Err(PublishError::new(
193 chain,
194 PublishErrorKind::Transport(SplitQueueError::CorruptState),
195 ));
196 }
197
198 storage.records[usize::from(head)] = Some(ChainRecord {
199 id,
200 state: ChainState::Available,
201 chain,
202 });
203 let ring_slot = storage.ring_slot(storage.available_index);
204 storage.available[ring_slot] = head;
205 storage.available_index = storage.available_index.wrapping_add(1);
206 storage.available_count += 1;
207 storage.active_chains += 1;
208
209 Ok(PublishedChain::new(
210 id,
211 if notifications {
212 NotificationHint::Notify
213 } else {
214 NotificationHint::Suppressed
215 },
216 ))
217 }
218
219 fn storage_mut(&mut self) -> Result<&mut RingStorage, QueueError<SplitQueueError>> {
220 self.storage.as_mut().ok_or(QueueError::NotReady)
221 }
222
223 fn advance_epoch(
224 &mut self,
225 next_epoch: QueueEpoch,
226 ) -> Result<Option<RingStorage>, QueueError<SplitQueueError>> {
227 if next_epoch <= self.state.epoch() {
228 return Err(QueueError::InvalidConfiguration(
229 QueueConfigError::NonIncreasingEpoch,
230 ));
231 }
232
233 self.current_epoch
234 .store(next_epoch.get(), Ordering::Release);
235 self.state = QueueState::unconfigured(self.state.max_size(), next_epoch);
236 self.available_notifications = true;
237 self.used_notifications = true;
238 Ok(self.storage.take())
239 }
240}
241
242impl QueuePort for SplitQueue {
243 fn state(&self) -> QueueState {
244 self.state
245 }
246}
247
248impl QueueControl for SplitQueue {
249 type Error = SplitQueueError;
250
251 fn configure(&mut self, size: QueueSize) -> Result<(), QueueError<Self::Error>> {
252 if self.state.ready() {
253 return Err(QueueError::Transport(SplitQueueError::Busy));
254 }
255 if size > self.state.max_size() {
256 return Err(QueueError::InvalidConfiguration(
257 QueueConfigError::SizeExceedsMaximum,
258 ));
259 }
260 if size.get() < self.max_chain_descriptors {
261 return Err(QueueError::Transport(
262 SplitQueueError::ChainLimitExceedsQueue,
263 ));
264 }
265 if self
266 .storage
267 .as_ref()
268 .is_some_and(|storage| storage.active_chains != 0)
269 {
270 return Err(QueueError::Transport(SplitQueueError::Busy));
271 }
272
273 let storage = RingStorage::new(size).map_err(QueueError::Transport)?;
274 self.state = QueueState::new(self.state.max_size(), Some(size), false, self.state.epoch())
275 .map_err(QueueError::InvalidConfiguration)?;
276 self.storage = Some(storage);
277 Ok(())
278 }
279
280 fn set_ready(&mut self, ready: bool) -> Result<(), QueueError<Self::Error>> {
281 self.state = QueueState::new(
282 self.state.max_size(),
283 self.state.size(),
284 ready,
285 self.state.epoch(),
286 )
287 .map_err(QueueError::InvalidConfiguration)?;
288 Ok(())
289 }
290}
291
292impl DriverQueue for SplitQueue {
293 type Chain = DriverChain;
294 type Reclaimed = ReclaimedChains;
295 type Error = SplitQueueError;
296
297 fn publish(
298 &mut self,
299 chain: Self::Chain,
300 ) -> Result<PublishedChain, PublishError<Self::Chain, Self::Error>> {
301 self.publish_inner(chain, false)
302 }
303
304 fn pop_used(&mut self) -> Result<Option<UsedChain<Self::Chain>>, QueueError<Self::Error>> {
305 if !self.state.ready() {
306 return Err(QueueError::NotReady);
307 }
308 let storage = self.storage_mut()?;
309 if storage.used_count == 0 {
310 return Ok(None);
311 }
312
313 let ring_slot = storage.ring_slot(storage.consumed_used);
314 let used = storage.used[ring_slot]
315 .take()
316 .ok_or(QueueError::Transport(SplitQueueError::CorruptState))?;
317 let Some(record) = storage.records[usize::from(used.head)].take() else {
318 storage.used[ring_slot] = Some(used);
319 return Err(QueueError::Transport(SplitQueueError::CorruptState));
320 };
321 if record.id != used.id || record.state != ChainState::Used {
322 storage.records[usize::from(used.head)] = Some(record);
323 storage.used[ring_slot] = Some(used);
324 return Err(QueueError::Transport(SplitQueueError::CorruptState));
325 }
326
327 storage.release_descriptors(&record.chain);
328 storage.consumed_used = storage.consumed_used.wrapping_add(1);
329 storage.used_count -= 1;
330 storage.active_chains -= 1;
331 Ok(Some(UsedChain::new(record.id, used.used, record.chain)))
332 }
333
334 fn disable_used_notifications(&mut self) -> Result<(), QueueError<Self::Error>> {
335 if !self.state.ready() {
336 return Err(QueueError::NotReady);
337 }
338 self.used_notifications = false;
339 Ok(())
340 }
341
342 fn enable_used_notifications(
343 &mut self,
344 ) -> Result<NotificationRecheck, QueueError<Self::Error>> {
345 if !self.state.ready() {
346 return Err(QueueError::NotReady);
347 }
348 self.used_notifications = true;
349 Ok(if self.storage_mut()?.used_count == 0 {
350 NotificationRecheck::Idle
351 } else {
352 NotificationRecheck::WorkPending
353 })
354 }
355
356 fn reset(
357 &mut self,
358 next_epoch: QueueEpoch,
359 ) -> Result<Self::Reclaimed, QueueError<Self::Error>> {
360 Ok(ReclaimedChains::new(self.advance_epoch(next_epoch)?))
361 }
362}
363
364impl DeviceQueue for SplitQueue {
365 type Chain = SplitDeviceChain;
366 type Error = SplitQueueError;
367
368 fn pop_available(&mut self) -> Result<Option<Self::Chain>, QueueError<Self::Error>> {
369 if !self.state.ready() {
370 return Err(QueueError::NotReady);
371 }
372 let epoch = self.state.epoch();
373 let max_descriptors = self.max_chain_descriptors;
374 let current_epoch = Rc::clone(&self.current_epoch);
375 let storage = self.storage_mut()?;
376 if storage.available_count == 0 {
377 return Ok(None);
378 }
379
380 let ring_slot = storage.ring_slot(storage.consumed_available);
381 let head = storage.available[ring_slot];
382 let record = storage.records[usize::from(head)]
383 .as_mut()
384 .ok_or(QueueError::Transport(SplitQueueError::CorruptState))?;
385 if record.id.epoch() != epoch || record.state != ChainState::Available {
386 return Err(QueueError::Transport(SplitQueueError::CorruptState));
387 }
388 record.state = ChainState::InFlight;
389 let chain = SplitDeviceChain::new(
390 record.id,
391 record.chain.data(),
392 max_descriptors,
393 current_epoch,
394 );
395 storage.consumed_available = storage.consumed_available.wrapping_add(1);
396 storage.available_count -= 1;
397 Ok(Some(chain))
398 }
399
400 fn complete(
401 &mut self,
402 chain: Self::Chain,
403 used: UsedLength,
404 ) -> Result<NotificationHint, QueueError<Self::Error>> {
405 let id = chain.id();
406 let current = self.state.epoch();
407 if id.epoch() != current {
408 return Err(QueueError::ResetRace {
409 operation: id.epoch(),
410 current,
411 });
412 }
413 if !self.state.ready() {
414 return Err(QueueError::NotReady);
415 }
416 let capacity = chain.writable_capacity();
417 if u64::from(used.get()) > capacity {
418 return Err(QueueError::UsedLengthExceeded { used, capacity });
419 }
420 let head = u16::try_from(id.token())
421 .map_err(|_| QueueError::Transport(SplitQueueError::CorruptState))?;
422 let notifications = self.used_notifications;
423 let storage = self.storage_mut()?;
424 let record = storage
425 .records
426 .get_mut(usize::from(head))
427 .and_then(Option::as_mut)
428 .ok_or(QueueError::Transport(SplitQueueError::CorruptState))?;
429 if record.id != id || record.state != ChainState::InFlight {
430 return Err(QueueError::Transport(SplitQueueError::CorruptState));
431 }
432 if storage.used_count == storage.size.get() {
433 return Err(QueueError::Transport(SplitQueueError::CorruptState));
434 }
435
436 record.state = ChainState::Used;
437 let ring_slot = storage.ring_slot(storage.used_index);
438 storage.used[ring_slot] = Some(UsedElement { head, id, used });
439 storage.used_index = storage.used_index.wrapping_add(1);
440 storage.used_count += 1;
441 Ok(if notifications {
442 NotificationHint::Notify
443 } else {
444 NotificationHint::Suppressed
445 })
446 }
447
448 fn disable_available_notifications(&mut self) -> Result<(), QueueError<Self::Error>> {
449 if !self.state.ready() {
450 return Err(QueueError::NotReady);
451 }
452 self.available_notifications = false;
453 Ok(())
454 }
455
456 fn enable_available_notifications(
457 &mut self,
458 ) -> Result<NotificationRecheck, QueueError<Self::Error>> {
459 if !self.state.ready() {
460 return Err(QueueError::NotReady);
461 }
462 self.available_notifications = true;
463 Ok(if self.storage_mut()?.available_count == 0 {
464 NotificationRecheck::Idle
465 } else {
466 NotificationRecheck::WorkPending
467 })
468 }
469
470 fn reset(&mut self, next_epoch: QueueEpoch) -> Result<(), QueueError<Self::Error>> {
471 drop(self.advance_epoch(next_epoch)?);
472 Ok(())
473 }
474}
475
476#[derive(Debug)]
477struct RingStorage {
478 size: QueueSize,
479 descriptors: Box<[Option<DescriptorOwner>]>,
480 records: Box<[Option<ChainRecord>]>,
481 available: Box<[u16]>,
482 used: Box<[Option<UsedElement>]>,
483 available_index: u16,
484 consumed_available: u16,
485 used_index: u16,
486 consumed_used: u16,
487 available_count: u16,
488 used_count: u16,
489 active_chains: u16,
490 free_descriptors: u16,
491 next_descriptor: u16,
492}
493
494impl RingStorage {
495 fn new(size: QueueSize) -> Result<Self, SplitQueueError> {
496 let len = usize::from(size.get());
497 Ok(Self {
498 size,
499 descriptors: boxed_with(len, |_| None)?,
500 records: boxed_with(len, |_| None)?,
501 available: boxed_with(len, |_| 0)?,
502 used: boxed_with(len, |_| None)?,
503 available_index: 0,
504 consumed_available: 0,
505 used_index: 0,
506 consumed_used: 0,
507 available_count: 0,
508 used_count: 0,
509 active_chains: 0,
510 free_descriptors: size.get(),
511 next_descriptor: 0,
512 })
513 }
514
515 fn ring_slot(&self, index: u16) -> usize {
516 usize::from(index & (self.size.get() - 1))
517 }
518
519 const fn counters(&self) -> RingCounters {
520 RingCounters {
521 available: self.available_index,
522 consumed_available: self.consumed_available,
523 used: self.used_index,
524 consumed_used: self.consumed_used,
525 }
526 }
527
528 fn allocate_descriptors(&mut self, chain: &mut DriverChain) -> Result<u16, SplitQueueError> {
529 let count = chain.descriptor_count();
530 if count > self.free_descriptors {
531 return Err(SplitQueueError::CorruptState);
532 }
533
534 let mut cursor = self.next_descriptor;
535 let mut allocated = 0_usize;
536 while allocated < usize::from(count) {
537 let mut scanned = 0_u16;
538 while self.descriptors[usize::from(cursor)].is_some() {
539 cursor = self.wrapping_descriptor_add(cursor, 1);
540 scanned += 1;
541 if scanned == self.size.get() {
542 self.rollback_allocation(chain, allocated);
543 return Err(SplitQueueError::CorruptState);
544 }
545 }
546 chain.slots_mut()[allocated] = cursor;
547 self.descriptors[usize::from(cursor)] = Some(DescriptorOwner {
548 head: 0,
549 local: allocated as u16,
550 });
551 allocated += 1;
552 cursor = self.wrapping_descriptor_add(cursor, 1);
553 }
554
555 let head = chain.queue_head_slot();
556 for slot in chain.slots() {
557 let owner = self.descriptors[usize::from(*slot)]
558 .as_mut()
559 .ok_or(SplitQueueError::CorruptState)?;
560 owner.head = head;
561 }
562 self.next_descriptor = cursor;
563 self.free_descriptors -= count;
564 Ok(head)
565 }
566
567 fn release_descriptors(&mut self, chain: &DriverChain) {
568 for slot in chain.slots() {
569 if let Some(owner) = self.descriptors[usize::from(*slot)].take() {
570 debug_assert_eq!(owner.head, chain.queue_head_slot());
571 debug_assert!(usize::from(owner.local) < chain.slots().len());
572 self.free_descriptors += 1;
573 }
574 }
575 }
576
577 fn rollback_allocation(&mut self, chain: &DriverChain, allocated: usize) {
578 for slot in &chain.slots()[..allocated] {
579 self.descriptors[usize::from(*slot)] = None;
580 }
581 }
582
583 const fn wrapping_descriptor_add(&self, index: u16, amount: u16) -> u16 {
584 index.wrapping_add(amount) & (self.size.get() - 1)
585 }
586}
587
588#[derive(Clone, Copy, Debug, PartialEq, Eq)]
589struct DescriptorOwner {
590 head: u16,
591 local: u16,
592}
593
594#[derive(Clone, Copy, Debug, PartialEq, Eq)]
595enum ChainState {
596 Available,
597 InFlight,
598 Used,
599}
600
601#[derive(Debug)]
602struct ChainRecord {
603 id: ChainId,
604 state: ChainState,
605 chain: DriverChain,
606}
607
608#[derive(Clone, Copy, Debug, PartialEq, Eq)]
609struct UsedElement {
610 head: u16,
611 id: ChainId,
612 used: UsedLength,
613}
614
615#[derive(Debug)]
617pub struct ReclaimedChains {
618 records: IntoIter<Option<ChainRecord>>,
619}
620
621impl ReclaimedChains {
622 fn new(storage: Option<RingStorage>) -> Self {
623 let records = storage
624 .map(|storage| storage.records.into_vec())
625 .unwrap_or_default()
626 .into_iter();
627 Self { records }
628 }
629}
630
631impl Iterator for ReclaimedChains {
632 type Item = ReclaimedChain<DriverChain>;
633
634 fn next(&mut self) -> Option<Self::Item> {
635 self.records
636 .by_ref()
637 .find_map(|record| record.map(|record| ReclaimedChain::new(record.id, record.chain)))
638 }
639}
640
641fn boxed_with<T>(
642 len: usize,
643 mut value: impl FnMut(usize) -> T,
644) -> Result<Box<[T]>, SplitQueueError> {
645 let mut values = Vec::new();
646 values
647 .try_reserve_exact(len)
648 .map_err(|_| SplitQueueError::AllocationFailed)?;
649 for index in 0..len {
650 values.push(value(index));
651 }
652 Ok(values.into_boxed_slice())
653}
654
655#[cfg(test)]
656mod tests {
657 extern crate std;
658
659 use std::vec;
660 use std::vec::Vec;
661
662 use virtio_accel_transport::{
663 ByteAccessError, ChainError, DeviceChain, ReadableBytes, WritableBytes,
664 };
665
666 use super::*;
667 use crate::{Descriptor, VIRTQ_DESC_F_INDIRECT, VIRTQ_DESC_F_NEXT, VIRTQ_DESC_F_WRITE};
668
669 fn ready_queue(size: u16, max_chain_descriptors: u16) -> SplitQueue {
670 let size = QueueSize::new(size).unwrap();
671 let mut queue = SplitQueue::new(size, max_chain_descriptors).unwrap();
672 QueueControl::configure(&mut queue, size).unwrap();
673 QueueControl::set_ready(&mut queue, true).unwrap();
674 queue
675 }
676
677 fn direct_chain(requests: &[&[u8]], responses: &[usize]) -> DriverChain {
678 let mut descriptors = Vec::new();
679 for request in requests {
680 descriptors.push(Descriptor::readable(request.to_vec()));
681 }
682 for response in responses {
683 descriptors.push(Descriptor::writable(vec![0; *response]));
684 }
685 DriverChain::direct(descriptors).unwrap()
686 }
687
688 fn assert_malformed(queue: &mut SplitQueue, chain: DriverChain, expected: MalformedChain) {
689 queue.inject_available(chain).unwrap();
690 let mut chain = DeviceQueue::pop_available(queue).unwrap().unwrap();
691 match chain.io() {
692 Err(ChainError::Malformed(actual)) => assert_eq!(actual, expected),
693 result => panic!("unexpected chain result: {result:?}"),
694 }
695 DeviceQueue::complete(queue, chain, UsedLength::new(0)).unwrap();
696 DriverQueue::pop_used(queue).unwrap().unwrap();
697 }
698
699 #[test]
700 fn vq_001_vq_003_segmented_chain_round_trip_does_not_coalesce_payloads() {
701 let mut queue = ready_queue(8, 4);
702 let published =
703 DriverQueue::publish(&mut queue, direct_chain(&[b"ab", b"cd"], &[2, 2])).unwrap();
704 assert_eq!(published.notification(), NotificationHint::Notify);
705
706 let mut device_chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
707 {
708 let (_, request, response) = device_chain.io().unwrap().into_parts();
709 let mut bytes = [0; 3];
710 request.read_at(1, &mut bytes).unwrap();
711 assert_eq!(&bytes, b"bcd");
712 response.write_at(0, b"wxyz").unwrap();
713 }
714 DeviceQueue::complete(&mut queue, device_chain, UsedLength::new(4)).unwrap();
715
716 let used = DriverQueue::pop_used(&mut queue).unwrap().unwrap();
717 assert_eq!(used.id(), published.id());
718 assert_eq!(used.used(), UsedLength::new(4));
719 let (_, _, chain) = used.into_parts();
720 let mut first = [0; 2];
721 let mut second = [0; 2];
722 chain.read_descriptor(2, 0, &mut first).unwrap();
723 chain.read_descriptor(3, 0, &mut second).unwrap();
724 assert_eq!(&first, b"wx");
725 assert_eq!(&second, b"yz");
726 assert_eq!(queue.free_descriptors(), Some(8));
727 }
728
729 #[test]
730 fn split_ring_and_descriptor_indices_wrap_naturally() {
731 let mut queue = ready_queue(8, 4);
732 queue.set_empty_ring_index(u16::MAX).unwrap();
733 queue.set_next_descriptor(7).unwrap();
734
735 let published = DriverQueue::publish(&mut queue, direct_chain(&[b"r"], &[1])).unwrap();
736 assert_eq!(published.id().token(), 7);
737 assert_eq!(
738 queue.storage.as_ref().unwrap().descriptors[0],
739 Some(DescriptorOwner { head: 7, local: 1 })
740 );
741 let chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
742 DeviceQueue::complete(&mut queue, chain, UsedLength::new(1)).unwrap();
743 DriverQueue::pop_used(&mut queue).unwrap().unwrap();
744
745 assert_eq!(
746 queue.ring_counters(),
747 Some(RingCounters {
748 available: 0,
749 consumed_available: 0,
750 used: 0,
751 consumed_used: 0,
752 })
753 );
754 }
755
756 #[test]
757 fn vq_016_descriptor_exhaustion_returns_the_unpublished_chain() {
758 let mut queue = ready_queue(4, 4);
759 DriverQueue::publish(&mut queue, direct_chain(&[b"a"], &[1])).unwrap();
760 DriverQueue::publish(&mut queue, direct_chain(&[b"b"], &[1])).unwrap();
761 assert_eq!(queue.free_descriptors(), Some(0));
762
763 let error = DriverQueue::publish(&mut queue, direct_chain(&[b"c"], &[1])).unwrap_err();
764 assert_eq!(error.kind(), &PublishErrorKind::InsufficientDescriptors);
765 let (chain, _) = error.into_parts();
766 assert_eq!(chain.descriptor_count(), 2);
767 assert_eq!(queue.free_descriptors(), Some(0));
768 }
769
770 #[test]
771 fn vq_004_vq_005_vq_006_malformed_topology_is_classified_boundedly() {
772 let mut queue = ready_queue(8, 4);
773
774 assert_malformed(
775 &mut queue,
776 DriverChain::raw(
777 vec![
778 Descriptor::raw(vec![1], VIRTQ_DESC_F_NEXT, 1),
779 Descriptor::raw(vec![0], VIRTQ_DESC_F_WRITE | VIRTQ_DESC_F_NEXT, 0),
780 ],
781 0,
782 )
783 .unwrap(),
784 MalformedChain::DescriptorLoop,
785 );
786 assert_malformed(
787 &mut queue,
788 DriverChain::raw(
789 vec![
790 Descriptor::raw(vec![1], VIRTQ_DESC_F_NEXT, 7),
791 Descriptor::writable(vec![0]),
792 ],
793 0,
794 )
795 .unwrap(),
796 MalformedChain::DescriptorIndex,
797 );
798 assert_malformed(
799 &mut queue,
800 DriverChain::raw(
801 vec![
802 Descriptor::raw(vec![1], 8, 0),
803 Descriptor::writable(vec![0]),
804 ],
805 0,
806 )
807 .unwrap(),
808 MalformedChain::DescriptorFlags,
809 );
810 assert_malformed(
811 &mut queue,
812 DriverChain::raw(
813 vec![
814 Descriptor::raw(vec![1], VIRTQ_DESC_F_INDIRECT, 0),
815 Descriptor::writable(vec![0]),
816 ],
817 0,
818 )
819 .unwrap(),
820 MalformedChain::IndirectUnsupported,
821 );
822 assert_malformed(
823 &mut queue,
824 DriverChain::raw(
825 vec![
826 Descriptor::raw(Vec::new(), VIRTQ_DESC_F_NEXT, 1),
827 Descriptor::writable(vec![0]),
828 ],
829 0,
830 )
831 .unwrap(),
832 MalformedChain::ZeroLength,
833 );
834 assert_malformed(
835 &mut queue,
836 DriverChain::raw(
837 vec![
838 Descriptor::raw(vec![0], VIRTQ_DESC_F_WRITE | VIRTQ_DESC_F_NEXT, 1),
839 Descriptor::readable(vec![1]),
840 ],
841 0,
842 )
843 .unwrap(),
844 MalformedChain::Direction,
845 );
846 assert_malformed(
847 &mut queue,
848 DriverChain::raw(
849 vec![
850 Descriptor::raw(vec![1], VIRTQ_DESC_F_NEXT, 1),
851 Descriptor::writable(vec![0]),
852 Descriptor::writable(vec![0]),
853 ],
854 0,
855 )
856 .unwrap(),
857 MalformedChain::DescriptorCount,
858 );
859 assert_malformed(
860 &mut queue,
861 DriverChain::raw(
862 vec![
863 Descriptor::raw(vec![0], VIRTQ_DESC_F_WRITE | VIRTQ_DESC_F_NEXT, 1),
864 Descriptor::writable(vec![0]),
865 ],
866 0,
867 )
868 .unwrap(),
869 MalformedChain::Direction,
870 );
871 assert_malformed(
872 &mut queue,
873 DriverChain::raw(
874 vec![
875 Descriptor::raw(vec![1], VIRTQ_DESC_F_NEXT, 1),
876 Descriptor::readable(vec![1]),
877 ],
878 0,
879 )
880 .unwrap(),
881 MalformedChain::Direction,
882 );
883 assert_malformed(
884 &mut queue,
885 DriverChain::raw(
886 vec![
887 Descriptor::unmapped(u64::MAX, VIRTQ_DESC_F_NEXT, 1),
888 Descriptor::writable(vec![0]),
889 ],
890 0,
891 )
892 .unwrap(),
893 MalformedChain::Address,
894 );
895 }
896
897 #[test]
898 fn vq_006_normal_publication_rejects_malformed_and_oversized_chains() {
899 let mut queue = ready_queue(8, 2);
900 let malformed = DriverChain::raw(
901 vec![
902 Descriptor::raw(Vec::new(), VIRTQ_DESC_F_NEXT, 1),
903 Descriptor::writable(vec![0]),
904 ],
905 0,
906 )
907 .unwrap();
908 let error = DriverQueue::publish(&mut queue, malformed).unwrap_err();
909 assert_eq!(
910 error.kind(),
911 &PublishErrorKind::Transport(SplitQueueError::MalformedDriverChain(
912 MalformedChain::ZeroLength
913 ))
914 );
915
916 let oversized = direct_chain(&[b"a", b"b"], &[1]);
917 let error = DriverQueue::publish(&mut queue, oversized).unwrap_err();
918 assert_eq!(error.kind(), &PublishErrorKind::InsufficientDescriptors);
919
920 let oversized = direct_chain(&[b"a", b"b"], &[1]);
921 queue.inject_available(oversized).unwrap();
922 let mut chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
923 assert!(matches!(
924 chain.io(),
925 Err(ChainError::Malformed(MalformedChain::DescriptorCount))
926 ));
927 }
928
929 #[test]
930 fn vq_014_completion_order_drives_used_ring_order() {
931 let mut queue = ready_queue(8, 4);
932 let first = DriverQueue::publish(&mut queue, direct_chain(&[b"a"], &[2]))
933 .unwrap()
934 .id();
935 let second = DriverQueue::publish(&mut queue, direct_chain(&[b"b"], &[2]))
936 .unwrap()
937 .id();
938 let first_chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
939 let second_chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
940
941 DeviceQueue::complete(&mut queue, second_chain, UsedLength::new(2)).unwrap();
942 DeviceQueue::complete(&mut queue, first_chain, UsedLength::new(1)).unwrap();
943 assert_eq!(
944 DriverQueue::pop_used(&mut queue).unwrap().unwrap().id(),
945 second
946 );
947 assert_eq!(
948 DriverQueue::pop_used(&mut queue).unwrap().unwrap().id(),
949 first
950 );
951 }
952
953 #[test]
954 fn vq_017_notification_enablement_rechecks_for_missed_work() {
955 let mut queue = ready_queue(8, 4);
956 DeviceQueue::disable_available_notifications(&mut queue).unwrap();
957 let published = DriverQueue::publish(&mut queue, direct_chain(&[b"a"], &[1])).unwrap();
958 assert_eq!(published.notification(), NotificationHint::Suppressed);
959 assert_eq!(
960 DeviceQueue::enable_available_notifications(&mut queue).unwrap(),
961 NotificationRecheck::WorkPending
962 );
963
964 let chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
965 assert_eq!(
966 DeviceQueue::enable_available_notifications(&mut queue).unwrap(),
967 NotificationRecheck::Idle
968 );
969 DriverQueue::disable_used_notifications(&mut queue).unwrap();
970 assert_eq!(
971 DeviceQueue::complete(&mut queue, chain, UsedLength::new(1)).unwrap(),
972 NotificationHint::Suppressed
973 );
974 assert_eq!(
975 DriverQueue::enable_used_notifications(&mut queue).unwrap(),
976 NotificationRecheck::WorkPending
977 );
978 DriverQueue::pop_used(&mut queue).unwrap().unwrap();
979 assert_eq!(
980 DriverQueue::enable_used_notifications(&mut queue).unwrap(),
981 NotificationRecheck::Idle
982 );
983 }
984
985 #[test]
986 fn vq_018_reset_invalidates_ports_before_reclaiming_driver_ownership() {
987 let mut queue = ready_queue(8, 4);
988 let published = DriverQueue::publish(&mut queue, direct_chain(&[b"a"], &[2])).unwrap();
989 let mut device_chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
990 let next_epoch = queue.state().epoch().checked_next().unwrap();
991
992 let mut reclaimed = {
993 let (_, _, response) = device_chain.io().unwrap().into_parts();
994 let reclaimed = DriverQueue::reset(&mut queue, next_epoch).unwrap();
995 assert_eq!(response.write_at(0, b"x"), Err(ByteAccessError::Reset));
996 reclaimed
997 };
998
999 let reclaimed_chain = reclaimed.next().unwrap();
1000 assert_eq!(reclaimed_chain.id(), published.id());
1001 assert!(reclaimed.next().is_none());
1002 assert!(matches!(
1003 device_chain.io(),
1004 Err(ChainError::ResetRace { chain, current })
1005 if chain == published.id().epoch() && current == next_epoch
1006 ));
1007 assert_eq!(
1008 DeviceQueue::complete(&mut queue, device_chain, UsedLength::new(0)),
1009 Err(QueueError::ResetRace {
1010 operation: published.id().epoch(),
1011 current: next_epoch,
1012 })
1013 );
1014 assert_eq!(
1015 queue.state(),
1016 QueueState::unconfigured(QueueSize::new(8).unwrap(), next_epoch)
1017 );
1018 }
1019
1020 #[test]
1021 fn vq_013_excessive_used_length_publishes_no_used_entry() {
1022 let mut queue = ready_queue(8, 4);
1023 DriverQueue::publish(&mut queue, direct_chain(&[b"a"], &[2])).unwrap();
1024 let chain = DeviceQueue::pop_available(&mut queue).unwrap().unwrap();
1025 assert_eq!(
1026 DeviceQueue::complete(&mut queue, chain, UsedLength::new(3)),
1027 Err(QueueError::UsedLengthExceeded {
1028 used: UsedLength::new(3),
1029 capacity: 2,
1030 })
1031 );
1032 assert_eq!(queue.ring_counters().unwrap().used, 0);
1033 assert!(DriverQueue::pop_used(&mut queue).unwrap().is_none());
1034 }
1035}