Skip to main content

virtio_accel_split_queue/
queue.rs

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/// Failure while constructing a split queue.
16#[derive(Clone, Copy, Debug, PartialEq, Eq)]
17pub enum SplitQueueInitError {
18    /// The chain limit must be at least two and no larger than the maximum queue size.
19    DescriptorLimit,
20}
21
22/// Concrete failure in the in-memory split-ring model.
23#[derive(Clone, Copy, Debug, PartialEq, Eq)]
24pub enum SplitQueueError {
25    /// Allocation of configured ring storage failed.
26    AllocationFailed,
27    /// Reconfiguration would discard chains still owned by the queue.
28    Busy,
29    /// The configured queue is smaller than the queue's descriptor-chain limit.
30    ChainLimitExceedsQueue,
31    /// A normal driver publication supplied a malformed chain.
32    MalformedDriverChain(MalformedChain),
33    /// A deterministic test hook supplied an invalid counter or descriptor index.
34    InvalidTestState,
35    /// Internal ring ownership does not match the published indices.
36    CorruptState,
37}
38
39/// Observable split-ring counters for deterministic wraparound tests.
40#[derive(Clone, Copy, Debug, PartialEq, Eq)]
41pub struct RingCounters {
42    /// Driver-published available index.
43    pub available: u16,
44    /// Device-consumed available index.
45    pub consumed_available: u16,
46    /// Device-published used index.
47    pub used: u16,
48    /// Driver-consumed used index.
49    pub consumed_used: u16,
50}
51
52/// Single-owner, bounded split-virtqueue reference model.
53///
54/// Configuration preallocates descriptor, available, used, and ownership tables. Publication,
55/// consumption, completion, notification control, and reset perform no allocation or payload copy.
56#[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    /// Construct an unconfigured queue.
68    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    /// Maximum flattened descriptor count accepted by normal publication.
87    pub const fn max_chain_descriptors(&self) -> u16 {
88        self.max_chain_descriptors
89    }
90
91    /// Inject a raw chain into the available ring, retaining malformed topology for device tests.
92    ///
93    /// This differs from [`DriverQueue::publish`] only by bypassing profile validation. Descriptor
94    /// table and available-ring capacity checks still apply.
95    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    /// Set all four ring indices while the configured queue owns no chains.
103    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    /// Select the descriptor index at which the next bounded allocation scan begins.
118    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    /// Return ring counters when storage is configured.
128    pub fn ring_counters(&self) -> Option<RingCounters> {
129        self.storage.as_ref().map(RingStorage::counters)
130    }
131
132    /// Return the number of currently free descriptor-table entries.
133    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/// Allocation-free iterator over chains reclaimed by a driver reset.
616#[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}