Skip to content

Commit 0453478

Browse files
authored
refactor(crab-usb): change queue logic (#79)
* Refactor USB host controller endpoint management and logging - Introduced a new structure to manage endpoint interfaces within the Device struct. - Enhanced endpoint setup logic to clear stale endpoints and manage interface-specific endpoints. - Implemented a method to find endpoint descriptors across all alternate settings. - Improved logging for transfer events, command completions, and port status changes with a budget mechanism to limit log output. - Updated control transfer handling to track and log control transfer states more effectively. - Modified the Device struct to maintain a map of claimed interfaces instead of a single current interface. - Adjusted endpoint retrieval methods to support multiple interfaces and their respective endpoints. - Enhanced error handling and logging for transfer events and reclaiming requests. * refactor: change logging level from info to debug for USB event handling - Updated logging statements in device and endpoint management to use debug level instead of info for better performance and reduced log verbosity. - Removed unused atomic log budget checks in host and transfer modules to streamline code. * refactor: streamline endpoint and event ring handling, enhance ISO transfer management * fix(event): handle transfer activity event when no other events are present * fix(endpoint): improve endpoint descriptor retrieval and streamline control request handling * fix(parser): optimize format type handling and clean up unused imports * fix(parser): remove unnecessary code and improve error handling in image conversion * fix(queue): replace UnsafeCell with Mutex for safer data handling and improve finished data management fix(endpoint): refactor ISO packet handling and streamline request completion logic fix(event): reduce event ring segments from 16 to 2 for optimized resource usage
1 parent 2a54eb6 commit 0453478

3 files changed

Lines changed: 139 additions & 53 deletions

File tree

usb-host/src/backend/kmod/queue.rs

Lines changed: 92 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,10 @@
11
use alloc::sync::Arc;
2-
use core::pin::Pin;
3-
use core::task::Context;
4-
use core::task::Poll;
52
use core::{
63
cell::UnsafeCell,
7-
sync::atomic::{AtomicBool, AtomicUsize, Ordering},
4+
hint::spin_loop,
5+
pin::Pin,
6+
sync::atomic::{AtomicBool, AtomicU8, AtomicUsize, Ordering},
7+
task::{Context, Poll},
88
};
99
use futures::task::AtomicWaker;
1010

@@ -35,40 +35,42 @@ impl<C> Clone for Finished<C> {
3535
}
3636

3737
pub struct FinishedInner<C> {
38-
data: UnsafeCell<BTreeMap<BusAddr, Arc<FinishedData<C>>>>,
38+
data: BTreeMap<BusAddr, Arc<FinishedData<C>>>,
3939
}
4040

41+
const SLOT_EMPTY: u8 = 0;
42+
const SLOT_WRITING: u8 = 1;
43+
const SLOT_READY: u8 = 2;
44+
const SLOT_READING: u8 = 3;
45+
4146
pub struct FinishedData<C> {
4247
taken: AtomicBool,
43-
finished: AtomicBool,
48+
state: AtomicU8,
4449
waker: AtomicWaker,
4550
data: UnsafeCell<Option<C>>,
4651
}
4752

4853
impl<C> FinishedData<C> {
4954
fn new() -> Self {
5055
Self {
51-
finished: AtomicBool::new(false),
56+
state: AtomicU8::new(SLOT_EMPTY),
5257
taken: AtomicBool::new(false),
5358
waker: AtomicWaker::new(),
5459
data: UnsafeCell::new(None),
5560
}
5661
}
5762
}
5863

59-
unsafe impl<C> Send for FinishedInner<C> {}
60-
unsafe impl<C> Sync for FinishedInner<C> {}
61-
unsafe impl<C> Send for FinishedData<C> {}
62-
unsafe impl<C> Sync for FinishedData<C> {}
64+
unsafe impl<C: Send> Send for FinishedData<C> {}
65+
unsafe impl<C: Send> Sync for FinishedData<C> {}
66+
67+
unsafe impl<C: Send> Send for FinishedInner<C> {}
68+
unsafe impl<C: Send> Sync for FinishedInner<C> {}
6369

6470
impl<C> FinishedInner<C> {
6571
fn clear_finished(&self, addr: BusAddr) {
66-
if let Some(data) = unsafe { &mut *self.data.get() }.get(&addr) {
67-
data.finished.store(false, Ordering::Release);
68-
data.taken.store(false, Ordering::Release);
69-
unsafe {
70-
(*data.data.get()).take();
71-
}
72+
if let Some(data) = self.data.get(&addr) {
73+
data.clear();
7274
}
7375
}
7476
}
@@ -81,9 +83,7 @@ impl<C> Finished<C> {
8183
data.insert(addr, Arc::new(FinishedData::new()));
8284
}
8385
Self {
84-
inner: Arc::new(FinishedInner {
85-
data: UnsafeCell::new(data),
86-
}),
86+
inner: Arc::new(FinishedInner { data }),
8787
}
8888
}
8989

@@ -92,13 +92,8 @@ impl<C> Finished<C> {
9292
}
9393

9494
pub fn set_finished(&self, addr: BusAddr, value: C) {
95-
let data = unsafe { &mut *self.inner.data.get() };
96-
if let Some(slot) = data.get_mut(&addr) {
97-
unsafe {
98-
*slot.data.get() = Some(value);
99-
}
100-
slot.finished.store(true, Ordering::Release);
101-
slot.waker.wake();
95+
if let Some(slot) = self.inner.data.get(&addr) {
96+
slot.set_finished(value);
10297
} else if take_queue_log_budget() {
10398
warn!(
10499
"usb queue: completion address {:#x} is not registered",
@@ -112,8 +107,7 @@ impl<C> Finished<C> {
112107
}
113108

114109
fn waiter(&self, addr: BusAddr) -> &FinishedData<C> {
115-
let data = unsafe { &mut *self.inner.data.get() };
116-
let slot = data.get(&addr).unwrap();
110+
let slot = self.inner.data.get(&addr).unwrap();
117111
if slot.taken.load(Ordering::Acquire) {
118112
panic!("waiter called after take_waiter");
119113
}
@@ -126,7 +120,7 @@ impl<C> Finished<C> {
126120
}
127121

128122
pub fn take_waiter(&self, addr: BusAddr) -> TWaiter<C> {
129-
let data = unsafe { &mut *self.inner.data.get() }.get(&addr).unwrap();
123+
let data = self.inner.data.get(&addr).unwrap();
130124
if data.taken.swap(true, Ordering::AcqRel) {
131125
panic!("take_waiter called multiple times for the same addr");
132126
}
@@ -161,10 +155,76 @@ impl<C> FinishedData<C> {
161155
self.waker.register(waker);
162156
}
163157

158+
fn clear(&self) {
159+
loop {
160+
match self.state.load(Ordering::Acquire) {
161+
SLOT_EMPTY => return,
162+
SLOT_READY => {
163+
if self
164+
.state
165+
.compare_exchange(
166+
SLOT_READY,
167+
SLOT_READING,
168+
Ordering::AcqRel,
169+
Ordering::Acquire,
170+
)
171+
.is_ok()
172+
{
173+
unsafe {
174+
(*self.data.get()).take();
175+
}
176+
self.state.store(SLOT_EMPTY, Ordering::Release);
177+
return;
178+
}
179+
}
180+
SLOT_WRITING | SLOT_READING => spin_loop(),
181+
_ => {
182+
self.state.store(SLOT_EMPTY, Ordering::Release);
183+
return;
184+
}
185+
}
186+
}
187+
}
188+
189+
pub fn set_finished(&self, value: C) {
190+
if self
191+
.state
192+
.compare_exchange(
193+
SLOT_EMPTY,
194+
SLOT_WRITING,
195+
Ordering::AcqRel,
196+
Ordering::Acquire,
197+
)
198+
.is_err()
199+
{
200+
if take_queue_log_budget() {
201+
warn!("usb queue: dropping duplicate completion for busy slot");
202+
}
203+
return;
204+
}
205+
206+
unsafe {
207+
*self.data.get() = Some(value);
208+
}
209+
self.state.store(SLOT_READY, Ordering::Release);
210+
self.waker.wake();
211+
}
212+
164213
pub fn get_finished(&self) -> Option<C> {
165-
if !self.finished.load(Ordering::Acquire) {
214+
if self
215+
.state
216+
.compare_exchange(
217+
SLOT_READY,
218+
SLOT_READING,
219+
Ordering::AcqRel,
220+
Ordering::Acquire,
221+
)
222+
.is_err()
223+
{
166224
return None;
167225
}
168-
unsafe { (*self.data.get()).take() }
226+
let value = unsafe { (*self.data.get()).take() };
227+
self.state.store(SLOT_EMPTY, Ordering::Release);
228+
value
169229
}
170230
}

usb-host/src/backend/kmod/xhci/endpoint.rs

Lines changed: 46 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,7 @@ enum SubmittedTdKind {
8585
#[derive(Clone, Copy)]
8686
struct IsoPacketTd {
8787
trb: TransferId,
88+
final_packet: bool,
8889
event: Option<TransferEvent>,
8990
actual: Option<usize>,
9091
}
@@ -365,21 +366,21 @@ impl Endpoint {
365366
&mut self,
366367
id: RequestId,
367368
request_id: EndpointRequestId,
368-
packets: &[IsoPacketTd],
369369
) -> Option<Result<TransferCompletion, TransferError>> {
370-
for (index, packet) in packets.iter().copied().enumerate() {
370+
let packet_count = self.iso_packet_count(request_id)?;
371+
for index in 0..packet_count {
371372
if self.iso_packet_done(request_id, index) {
372373
continue;
373374
}
375+
let (packet_trb, requested) = self.iso_packet_info(request_id, index)?;
374376

375-
let Some(event) = self.ring.get_finished(packet.trb.0) else {
377+
let Some(event) = self.ring.get_finished(packet_trb.0) else {
376378
continue;
377379
};
378-
let requested = self.iso_requested_length(request_id, index)?;
379380
let actual = match iso_packet_actual_length(requested, event) {
380381
Ok(actual) => actual,
381382
Err(err) => {
382-
let cleanup_result = self.complete_request(request_id, packet.trb, event);
383+
let cleanup_result = self.complete_request(request_id, packet_trb, event);
383384
let result = match cleanup_result {
384385
Ok(_) => Err(err),
385386
Err(cleanup_err) => Err(cleanup_err),
@@ -389,17 +390,27 @@ impl Endpoint {
389390
};
390391

391392
let fatal = iso_packet_is_fatal(event);
392-
let all_completed = self.record_iso_packet(request_id, index, event, actual)?;
393-
if fatal || all_completed {
393+
let should_complete =
394+
self.record_iso_packet(request_id, index, event, actual, fatal)?;
395+
if should_complete {
394396
return Some(
395-
self.complete_request(request_id, packet.trb, event)
397+
self.complete_request(request_id, packet_trb, event)
396398
.map(|transfer| transfer_to_completion(id, transfer)),
397399
);
398400
}
399401
}
400402
None
401403
}
402404

405+
fn iso_packet_count(&self, request_id: EndpointRequestId) -> Option<usize> {
406+
self.inflight
407+
.get(&request_id)
408+
.and_then(|submitted| match &submitted.kind {
409+
SubmittedTdKind::Iso { packets } => Some(packets.len()),
410+
_ => None,
411+
})
412+
}
413+
403414
fn iso_packet_done(&self, request_id: EndpointRequestId, index: usize) -> bool {
404415
self.inflight
405416
.get(&request_id)
@@ -412,13 +423,23 @@ impl Endpoint {
412423
.unwrap_or(true)
413424
}
414425

415-
fn iso_requested_length(&self, request_id: EndpointRequestId, index: usize) -> Option<usize> {
416-
self.inflight
417-
.get(&request_id)
418-
.and_then(|submitted| match &submitted.transfer.kind {
419-
TransferKind::Isochronous { packet_lengths } => packet_lengths.get(index).copied(),
420-
_ => None,
421-
})
426+
fn iso_packet_info(
427+
&self,
428+
request_id: EndpointRequestId,
429+
index: usize,
430+
) -> Option<(TransferId, usize)> {
431+
let submitted = self.inflight.get(&request_id)?;
432+
let SubmittedTdKind::Iso { packets } = &submitted.kind else {
433+
return None;
434+
};
435+
let packet = packets.get(index)?;
436+
let requested = match &submitted.transfer.kind {
437+
TransferKind::Isochronous { packet_lengths } => {
438+
packet_lengths.get(index).copied().unwrap_or(0)
439+
}
440+
_ => return None,
441+
};
442+
Some((packet.trb, requested))
422443
}
423444

424445
fn record_iso_packet(
@@ -427,21 +448,25 @@ impl Endpoint {
427448
index: usize,
428449
event: TransferEvent,
429450
actual: usize,
451+
fatal: bool,
430452
) -> Option<bool> {
431453
let submitted = self.inflight.get_mut(&request_id)?;
432454
let SubmittedTdKind::Iso { packets } = &mut submitted.kind else {
433455
return None;
434456
};
435-
for packet in packets.iter_mut().take(index) {
436-
if packet.actual.is_none() {
437-
packet.actual = Some(0);
457+
let final_packet = packets.get(index).is_some_and(|packet| packet.final_packet);
458+
if final_packet || fatal {
459+
for packet in packets.iter_mut().take(index) {
460+
if packet.actual.is_none() {
461+
packet.actual = Some(0);
462+
}
438463
}
439464
}
440465
if let Some(packet) = packets.get_mut(index) {
441466
packet.event = Some(event);
442467
packet.actual = Some(actual);
443468
}
444-
Some(packets.iter().all(|packet| packet.actual.is_some()))
469+
Some(final_packet || fatal || packets.iter().all(|packet| packet.actual.is_some()))
445470
}
446471

447472
fn enque_trb(&mut self, trb: transfer::Allowed) -> TransferId {
@@ -473,6 +498,7 @@ impl Endpoint {
473498
);
474499
packets.push(IsoPacketTd {
475500
trb,
501+
final_packet: last_packet,
476502
event: None,
477503
actual: None,
478504
});
@@ -746,7 +772,7 @@ impl EndpointOp for Endpoint {
746772
SubmittedTdKind::Control(control_td) => {
747773
self.reclaim_control_request(id, request_id, control_td)
748774
}
749-
SubmittedTdKind::Iso { packets } => self.reclaim_iso_request(id, request_id, &packets),
775+
SubmittedTdKind::Iso { .. } => self.reclaim_iso_request(id, request_id),
750776
}
751777
}
752778

usb-host/src/backend/kmod/xhci/event.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ pub struct EventRing {
2525
unsafe impl Send for EventRing {}
2626
unsafe impl Sync for EventRing {}
2727

28-
const EVENT_RING_SEGMENTS: usize = 16;
28+
const EVENT_RING_SEGMENTS: usize = 2;
2929

3030
impl EventRing {
3131
pub fn new(max_segments: usize, dma: &Kernel) -> Result<Self> {

0 commit comments

Comments
 (0)