1use fdf_sys::*;
8use libasync_dispatcher::{
9 AsAsyncDispatcherRef, AsyncDispatcher, AsyncDispatcherRef, GetAsyncDispatcher, JoinHandle,
10 OnDispatcher, Task,
11};
12
13use core::cell::RefCell;
14use core::ffi;
15use core::marker::PhantomData;
16use core::mem::ManuallyDrop;
17use core::ptr::{NonNull, null_mut};
18
19use zx::Status;
20
21use crate::shutdown_observer::ShutdownObserver;
22
23pub use fdf_sys::fdf_dispatcher_t;
24
25pub trait ShutdownObserverFn: FnOnce(DriverDispatcherRef<'_>) + Send + 'static {}
27impl<T> ShutdownObserverFn for T where T: FnOnce(DriverDispatcherRef<'_>) + Send + 'static {}
28
29#[derive(Default)]
31pub struct DispatcherBuilder {
32 #[doc(hidden)]
33 pub options: u32,
34 #[doc(hidden)]
35 pub name: String,
36 #[doc(hidden)]
37 pub scheduler_role: String,
38 #[doc(hidden)]
39 pub shutdown_observer: Option<Box<dyn ShutdownObserverFn>>,
40}
41
42impl DispatcherBuilder {
43 pub(crate) const UNSYNCHRONIZED: u32 = fdf_sys::FDF_DISPATCHER_OPTION_UNSYNCHRONIZED;
45 pub(crate) const ALLOW_THREAD_BLOCKING: u32 = fdf_sys::FDF_DISPATCHER_OPTION_ALLOW_SYNC_CALLS;
47 pub(crate) const NO_THREAD_MIGRATION: u32 = fdf_sys::FDF_DISPATCHER_OPTION_NO_THREAD_MIGRATION;
49
50 pub fn new() -> Self {
54 Self::default()
55 }
56
57 pub fn unsynchronized(mut self) -> Self {
63 assert!(
64 !self.allows_thread_blocking(),
65 "you may not create an unsynchronized dispatcher that allows synchronous calls"
66 );
67 self.options |= Self::UNSYNCHRONIZED;
68 self
69 }
70
71 pub fn is_unsynchronized(&self) -> bool {
73 (self.options & Self::UNSYNCHRONIZED) == Self::UNSYNCHRONIZED
74 }
75
76 pub fn allow_thread_blocking(mut self) -> Self {
82 assert!(
83 !self.is_unsynchronized(),
84 "you may not create an unsynchronized dispatcher that allows synchronous calls"
85 );
86 self.options |= Self::ALLOW_THREAD_BLOCKING;
87 self
88 }
89
90 pub fn allows_thread_blocking(&self) -> bool {
92 (self.options & Self::ALLOW_THREAD_BLOCKING) == Self::ALLOW_THREAD_BLOCKING
93 }
94
95 pub fn no_thread_migration(mut self) -> Self {
102 self.options |= Self::NO_THREAD_MIGRATION;
103 self
104 }
105
106 pub fn allows_thread_migration(&self) -> bool {
108 (self.options & Self::NO_THREAD_MIGRATION) == 0
109 }
110
111 pub fn name(mut self, name: &str) -> Self {
114 self.name = name.to_string();
115 self
116 }
117
118 pub fn scheduler_role(mut self, role: &str) -> Self {
122 self.scheduler_role = role.to_string();
123 self
124 }
125
126 pub fn shutdown_observer<F: ShutdownObserverFn>(mut self, shutdown_observer: F) -> Self {
128 self.shutdown_observer = Some(Box::new(shutdown_observer));
129 self
130 }
131
132 pub fn create(self) -> Result<Dispatcher, Status> {
137 let mut out_dispatcher = null_mut();
138 let options = self.options;
139 let name = self.name.as_ptr() as *mut ffi::c_char;
140 let name_len = self.name.len();
141 let scheduler_role = self.scheduler_role.as_ptr() as *mut ffi::c_char;
142 let scheduler_role_len = self.scheduler_role.len();
143 let observer =
144 ShutdownObserver::new(self.shutdown_observer.unwrap_or_else(|| Box::new(|_| {})))
145 .into_ptr();
146 Status::ok(unsafe {
150 fdf_dispatcher_create(
151 options,
152 name,
153 name_len,
154 scheduler_role,
155 scheduler_role_len,
156 observer,
157 &mut out_dispatcher,
158 )
159 })?;
160 Ok(Dispatcher(unsafe { NonNull::new_unchecked(out_dispatcher) }))
163 }
164
165 pub fn create_released(self) -> Result<AutoReleaseDispatcher, Status> {
169 self.create().map(Dispatcher::release)
170 }
171}
172
173#[derive(Debug)]
175pub struct Dispatcher(pub(crate) NonNull<fdf_dispatcher_t>);
176
177unsafe impl Send for Dispatcher {}
179unsafe impl Sync for Dispatcher {}
180thread_local! {
181 pub(crate) static OVERRIDE_DISPATCHER: RefCell<Option<NonNull<fdf_dispatcher_t>>> = const { RefCell::new(None) };
182}
183
184impl Dispatcher {
185 pub unsafe fn from_raw(handle: NonNull<fdf_dispatcher_t>) -> Self {
193 Self(handle)
194 }
195
196 fn get_raw_flags(&self) -> u32 {
197 unsafe { fdf_dispatcher_get_options(self.0.as_ptr()) }
199 }
200
201 pub fn is_unsynchronized(&self) -> bool {
203 (self.get_raw_flags() & DispatcherBuilder::UNSYNCHRONIZED) != 0
204 }
205
206 pub fn allows_thread_blocking(&self) -> bool {
208 (self.get_raw_flags() & DispatcherBuilder::ALLOW_THREAD_BLOCKING) != 0
209 }
210
211 pub fn allows_thread_migration(&self) -> bool {
214 (self.get_raw_flags() & DispatcherBuilder::NO_THREAD_MIGRATION) == 0
215 }
216
217 pub fn is_current_dispatcher(&self) -> bool {
219 self.0.as_ptr() == unsafe { fdf_dispatcher_get_current_dispatcher() }
222 }
223
224 pub fn release(self) -> AutoReleaseDispatcher {
229 AutoReleaseDispatcher { dispatcher: ManuallyDrop::new(self) }
230 }
231
232 pub fn as_dispatcher_ref(&self) -> DriverDispatcherRef<'_> {
235 DriverDispatcherRef(ManuallyDrop::new(Dispatcher(self.0)), PhantomData)
236 }
237}
238
239impl AsAsyncDispatcherRef for Dispatcher {
240 fn as_async_dispatcher_ref(&self) -> AsyncDispatcherRef<'_> {
241 let async_dispatcher =
242 NonNull::new(unsafe { fdf_dispatcher_get_async_dispatcher(self.0.as_ptr()) })
243 .expect("No async dispatcher on driver dispatcher");
244 unsafe { AsyncDispatcherRef::from_raw(async_dispatcher) }
245 }
246}
247
248impl Drop for Dispatcher {
249 fn drop(&mut self) {
250 unsafe { fdf_dispatcher_shutdown_async(self.0.as_mut()) }
253 }
254}
255
256#[derive(Debug)]
266pub struct AutoReleaseDispatcher {
267 dispatcher: ManuallyDrop<Dispatcher>,
268}
269
270impl AutoReleaseDispatcher {
271 pub unsafe fn from_raw(dispatcher: NonNull<fdf_dispatcher_t>) -> Self {
279 let dispatcher = ManuallyDrop::new(Dispatcher(dispatcher));
280 Self { dispatcher }
281 }
282
283 pub fn as_async_dispatcher(&self) -> AsyncDispatcher {
287 AsyncDispatcher::new(self)
288 }
289
290 pub fn as_dispatcher_ref(&self) -> DriverDispatcherRef<'_> {
293 DriverDispatcherRef(ManuallyDrop::new(Dispatcher(self.dispatcher.0)), PhantomData)
294 }
295
296 pub fn always_on_dispatcher(&self) -> AutoReleaseDispatcher {
298 let dispatcher_ref = unsafe { DriverDispatcherRef::from_raw(self.dispatcher.0) };
301 let dispatcher = unsafe { Dispatcher::from_raw(dispatcher_ref.always_on_dispatcher().0.0) };
306 Self { dispatcher: ManuallyDrop::new(dispatcher) }
307 }
308}
309
310impl AsAsyncDispatcherRef for AutoReleaseDispatcher {
311 fn as_async_dispatcher_ref(&self) -> AsyncDispatcherRef<'_> {
312 self.dispatcher.as_async_dispatcher_ref()
313 }
314}
315
316impl From<Dispatcher> for AutoReleaseDispatcher {
317 fn from(dispatcher: Dispatcher) -> Self {
318 Self { dispatcher: ManuallyDrop::new(dispatcher) }
319 }
320}
321
322#[derive(Debug)]
326pub struct DriverDispatcherRef<'a>(ManuallyDrop<Dispatcher>, PhantomData<&'a Dispatcher>);
327
328impl<'a> DriverDispatcherRef<'a> {
329 pub unsafe fn from_raw(handle: NonNull<fdf_dispatcher_t>) -> Self {
336 Self(ManuallyDrop::new(unsafe { Dispatcher::from_raw(handle) }), PhantomData)
338 }
339
340 pub fn from_async_dispatcher(dispatcher: AsyncDispatcherRef<'a>) -> Self {
347 let handle = NonNull::new(unsafe {
348 fdf_dispatcher_downcast_async_dispatcher(dispatcher.inner().as_ptr())
349 })
350 .unwrap();
351 unsafe { Self::from_raw(handle) }
352 }
353
354 pub unsafe fn as_raw(&mut self) -> *mut fdf_dispatcher_t {
360 unsafe { self.0.0.as_mut() }
361 }
362
363 pub fn always_on_dispatcher(&self) -> DriverDispatcherRef<'a> {
366 let ptr = unsafe { fdf_dispatcher_get_always_on_dispatcher(self.0.0.as_ptr()) };
368 DriverDispatcherRef(
369 ManuallyDrop::new(Dispatcher(NonNull::new(ptr).expect("Always-on dispatcher is NULL"))),
370 PhantomData,
371 )
372 }
373
374 pub fn register_wake_vector(
379 &self,
380 handle: &impl zx::AsHandleRef,
381 signals: zx::Signals,
382 ) -> Result<WakeVectorRegistration, Status> {
383 let raw_handle = handle.as_handle_ref().raw_handle();
384 Status::ok(unsafe {
387 fdf_sys::fdf_dispatcher_register_wake_vector(
388 self.0.0.as_ptr(),
389 raw_handle,
390 signals.bits(),
391 )
392 })?;
393 Ok(WakeVectorRegistration {
394 dispatcher: Some(AsyncDispatcher::new(self)),
395 handle: raw_handle,
396 signals,
397 })
398 }
399}
400
401#[derive(Debug)]
408pub struct WakeVectorRegistration {
409 dispatcher: Option<AsyncDispatcher>,
411 handle: zx::sys::zx_handle_t,
412 signals: zx::Signals,
413}
414
415impl WakeVectorRegistration {
416 pub fn unregister(mut self) -> Result<(), Status> {
422 self.unregister_inner()
423 }
424
425 fn unregister_inner(&mut self) -> Result<(), Status> {
426 let Some(dispatcher) = self.dispatcher.take() else { return Ok(()) };
427 let dispatcher_ref =
428 DriverDispatcherRef::from_async_dispatcher(dispatcher.as_async_dispatcher_ref());
429 Status::ok(unsafe {
432 fdf_sys::fdf_dispatcher_unregister_wake_vector(
433 dispatcher_ref.0.0.as_ptr(),
434 self.handle,
435 self.signals.bits(),
436 )
437 })
438 }
439}
440
441impl Drop for WakeVectorRegistration {
442 fn drop(&mut self) {
443 let _: Result<(), Status> = self.unregister_inner();
446 }
447}
448
449struct AddSendFuture<T>(T);
457
458impl<T: Future> Future for AddSendFuture<T> {
459 type Output = T::Output;
460
461 fn poll(
462 self: std::pin::Pin<&mut Self>,
463 cx: &mut std::task::Context<'_>,
464 ) -> std::task::Poll<Self::Output> {
465 let fut = unsafe { self.map_unchecked_mut(|fut| &mut fut.0) };
467 fut.poll(cx)
468 }
469}
470
471unsafe impl<T> Send for AddSendFuture<T> {}
475
476pub trait OnDriverDispatcher: OnDispatcher {
479 fn spawn_local(&self, future: impl Future<Output = ()> + 'static) -> JoinHandle<()>
492 where
493 Self: 'static,
494 {
495 self.compute_local(future).detach_on_drop()
496 }
497
498 fn compute_local<T: Send + 'static>(&self, future: impl Future<Output = T> + 'static) -> Task<T>
513 where
514 Self: 'static,
515 {
516 let Some(dispatcher) = self.try_get_async_dispatcher() else {
517 return Task::new_failed(Status::BAD_STATE);
518 };
519 let dispatcher =
520 DriverDispatcherRef::from_async_dispatcher(dispatcher.as_async_dispatcher_ref());
521 if dispatcher.0.is_current_dispatcher() && !dispatcher.0.allows_thread_migration() {
522 OnDispatcher::compute(self, AddSendFuture(future))
523 } else {
524 Task::new_failed(Status::BAD_STATE)
525 }
526 }
527
528 fn register_wake_vector(
530 &self,
531 handle: &impl zx::AsHandleRef,
532 signals: zx::Signals,
533 ) -> Result<WakeVectorRegistration, Status> {
534 let dispatcher = self.try_get_async_dispatcher().ok_or(Status::BAD_STATE)?;
535 let dispatcher_ref =
536 DriverDispatcherRef::from_async_dispatcher(dispatcher.as_async_dispatcher_ref());
537 dispatcher_ref.register_wake_vector(handle, signals)
538 }
539}
540
541impl<'a> AsAsyncDispatcherRef for DriverDispatcherRef<'a> {
542 fn as_async_dispatcher_ref(&self) -> AsyncDispatcherRef<'_> {
543 self.0.as_async_dispatcher_ref()
544 }
545}
546
547impl<'a> Clone for DriverDispatcherRef<'a> {
548 fn clone(&self) -> Self {
549 Self(ManuallyDrop::new(Dispatcher(self.0.0)), PhantomData)
550 }
551}
552
553impl<'a> core::ops::Deref for DriverDispatcherRef<'a> {
554 type Target = Dispatcher;
555 fn deref(&self) -> &Self::Target {
556 &self.0
557 }
558}
559
560impl<'a> core::ops::DerefMut for DriverDispatcherRef<'a> {
561 fn deref_mut(&mut self) -> &mut Self::Target {
562 &mut self.0
563 }
564}
565
566impl<T> OnDriverDispatcher for T where T: AsAsyncDispatcherRef + Clone {}
569
570#[derive(Clone, Copy, Debug, Default, PartialEq)]
572pub struct CurrentDispatcher;
573
574impl GetAsyncDispatcher for CurrentDispatcher {
575 fn try_get_async_dispatcher(&self) -> Option<AsyncDispatcher> {
576 OVERRIDE_DISPATCHER
577 .with(|global| *global.borrow())
578 .or_else(|| {
579 NonNull::new(unsafe { fdf_dispatcher_get_current_dispatcher() })
581 })
582 .map(|dispatcher| {
583 let async_dispatcher = NonNull::new(unsafe {
589 fdf_dispatcher_get_async_dispatcher(dispatcher.as_ptr())
590 })
591 .expect("No async dispatcher on driver dispatcher");
592 AsyncDispatcher::new(&unsafe { AsyncDispatcherRef::from_raw(async_dispatcher) })
593 })
594 }
595}
596
597impl OnDriverDispatcher for CurrentDispatcher {}
598
599#[cfg(test)]
600mod tests {
601 use super::*;
602
603 use std::sync::{Once, mpsc};
604
605 use futures::channel::mpsc as async_mpsc;
606 use futures::{SinkExt, StreamExt};
607 use zx::sys::ZX_OK;
608
609 use core::ffi::{c_char, c_void};
610 use core::ptr::null_mut;
611
612 static GLOBAL_DRIVER_ENV: Once = Once::new();
613 const NO_SYNC_CALLS_ROLE: &str = "no sync calls role";
614
615 pub fn ensure_driver_env() {
616 GLOBAL_DRIVER_ENV.call_once(|| {
617 unsafe {
620 assert_eq!(fdf_env_start(0), ZX_OK);
621 assert_eq!(
622 fdf_env_set_scheduler_role_opts(
623 NO_SYNC_CALLS_ROLE.as_ptr() as *const c_char,
624 NO_SYNC_CALLS_ROLE.len(),
625 FDF_SCHEDULER_ROLE_OPTION_NO_SYNC_CALLS
626 ),
627 ZX_OK
628 );
629 }
630 });
631 }
632 pub fn with_raw_dispatcher<T>(name: &str, p: impl FnOnce(AsyncDispatcher) -> T) -> T {
633 with_raw_dispatcher_flags(name, DispatcherBuilder::ALLOW_THREAD_BLOCKING, "", p)
634 }
635
636 pub(crate) fn with_raw_dispatcher_flags<T>(
637 name: &str,
638 flags: u32,
639 scheduler_role: &str,
640 p: impl FnOnce(AsyncDispatcher) -> T,
641 ) -> T {
642 ensure_driver_env();
643
644 let (shutdown_tx, shutdown_rx) = mpsc::channel();
645 let mut dispatcher = null_mut();
646 let mut observer = ShutdownObserver::new(move |dispatcher| {
647 assert!(!unsafe { fdf_env_dispatcher_has_queued_tasks(dispatcher.0.0.as_ptr()) });
650 shutdown_tx.send(()).unwrap();
651 })
652 .into_ptr();
653 let driver_ptr = &mut observer as *mut _ as *mut c_void;
654 let res = unsafe {
659 fdf_env_dispatcher_create_with_owner(
660 driver_ptr,
661 flags,
662 name.as_ptr() as *const c_char,
663 name.len(),
664 scheduler_role.as_ptr() as *const c_char,
665 scheduler_role.len(),
666 observer,
667 &mut dispatcher,
668 )
669 };
670 assert_eq!(res, ZX_OK);
671 let dispatcher = Dispatcher(NonNull::new(dispatcher).unwrap());
672
673 let res = p(AsyncDispatcher::new(&dispatcher));
674
675 drop(dispatcher);
676 shutdown_rx.recv().unwrap();
677
678 res
679 }
680
681 #[test]
682 fn start_test_dispatcher() {
683 with_raw_dispatcher("testing", |dispatcher| {
684 println!("hello {dispatcher:?}");
685 })
686 }
687
688 #[test]
689 fn post_task_on_dispatcher() {
690 with_raw_dispatcher("testing task", |dispatcher| {
691 let (tx, rx) = mpsc::channel();
692 dispatcher
693 .post_task_sync(move |status| {
694 assert_eq!(status, Ok(()));
695 tx.send(status).unwrap();
696 })
697 .unwrap();
698 assert_eq!(rx.recv().unwrap(), Ok(()));
699 });
700 }
701
702 #[test]
703 fn post_task_on_subdispatcher() {
704 let (shutdown_tx, shutdown_rx) = mpsc::channel();
705 with_raw_dispatcher("testing task top level", move |dispatcher| {
706 let (tx, rx) = mpsc::channel();
707 let (inner_tx, inner_rx) = mpsc::channel();
708 dispatcher
709 .post_task_sync(move |status| {
710 assert_eq!(status, Ok(()));
711 let inner = DispatcherBuilder::new()
712 .name("testing task second level")
713 .scheduler_role("")
714 .allow_thread_blocking()
715 .shutdown_observer(move |_dispatcher| {
716 println!("shutdown observer called");
717 shutdown_tx.send(1).unwrap();
718 })
719 .create()
720 .unwrap();
721 inner
722 .post_task_sync(move |status| {
723 assert_eq!(status, Ok(()));
724 tx.send(status).unwrap();
725 })
726 .unwrap();
727 inner_tx.send(inner).unwrap();
731 })
732 .unwrap();
733 assert_eq!(rx.recv().unwrap(), Ok(()));
734 inner_rx.recv().unwrap();
735 });
736 assert_eq!(shutdown_rx.recv().unwrap(), 1);
737 }
738
739 #[test]
740 fn spawn_local_fails_on_normal_dispatcher() {
741 let (shutdown_tx, shutdown_rx) = mpsc::channel();
742 with_raw_dispatcher("spawn local failures", move |dispatcher| {
743 let inside_dispatcher = dispatcher.clone();
744 dispatcher.spawn(async move {
745 assert_eq!(
746 inside_dispatcher.spawn_local(futures::future::ready(())).await.unwrap_err(),
747 Status::BAD_STATE
748 );
749 assert_eq!(
750 inside_dispatcher.compute_local(futures::future::ready(())).await.unwrap_err(),
751 Status::BAD_STATE
752 );
753 shutdown_tx.send(()).unwrap();
754 });
755 shutdown_rx.recv().unwrap();
756 });
757 }
758
759 #[test]
760 #[ignore = "Pending resolution of b/488397193"]
761 fn spawn_local_succeeds_on_no_thread_migration_dispatcher() {
762 let (tx, rx) = mpsc::channel();
763 with_raw_dispatcher_flags(
764 "spawn local success",
765 FDF_DISPATCHER_OPTION_NO_THREAD_MIGRATION,
766 NO_SYNC_CALLS_ROLE,
767 move |dispatcher| {
768 let inside_dispatcher = dispatcher.clone();
769 dispatcher.spawn(async move {
770 let tx_clone = tx.clone();
771 inside_dispatcher.spawn_local(async move {
772 tx_clone.send(()).unwrap();
773 });
774 inside_dispatcher
775 .compute_local(async move {
776 tx.send(()).unwrap();
777 })
778 .await
779 .unwrap();
780 });
781 rx.recv().unwrap();
783 rx.recv().unwrap();
784 },
785 );
786 }
787
788 #[test]
789 #[ignore = "Pending resolution of b/488397193"]
790 fn spawn_local_fails_on_no_thread_migration_dispatcher_from_different_thread() {
791 with_raw_dispatcher_flags(
792 "spawn local success",
793 FDF_DISPATCHER_OPTION_NO_THREAD_MIGRATION,
794 NO_SYNC_CALLS_ROLE,
795 move |dispatcher| {
796 let mut executor = fuchsia_async::LocalExecutor::default();
797 executor.run_singlethreaded(async {
798 assert_eq!(
801 dispatcher.spawn_local(futures::future::ready(())).await.unwrap_err(),
802 Status::BAD_STATE
803 );
804 assert_eq!(
805 dispatcher.compute_local(futures::future::ready(())).await.unwrap_err(),
806 Status::BAD_STATE
807 );
808 });
809 },
810 );
811 }
812
813 async fn ping(mut tx: async_mpsc::Sender<u8>, mut rx: async_mpsc::Receiver<u8>) {
814 println!("starting ping!");
815 tx.send(0).await.unwrap();
816 while let Some(next) = rx.next().await {
817 println!("ping! {next}");
818 tx.send(next + 1).await.unwrap();
819 }
820 }
821
822 async fn pong(
823 fin_tx: std::sync::mpsc::Sender<()>,
824 mut tx: async_mpsc::Sender<u8>,
825 mut rx: async_mpsc::Receiver<u8>,
826 ) {
827 println!("starting pong!");
828 while let Some(next) = rx.next().await {
829 println!("pong! {next}");
830 if next > 10 {
831 println!("bye!");
832 break;
833 }
834 tx.send(next + 1).await.unwrap();
835 }
836 fin_tx.send(()).unwrap();
837 }
838
839 #[test]
840 fn async_ping_pong() {
841 with_raw_dispatcher("async ping pong", |dispatcher| {
842 let (fin_tx, fin_rx) = mpsc::channel();
843 let (ping_tx, pong_rx) = async_mpsc::channel(10);
844 let (pong_tx, ping_rx) = async_mpsc::channel(10);
845 dispatcher.spawn(ping(ping_tx, ping_rx));
846 dispatcher.spawn(pong(fin_tx, pong_tx, pong_rx));
847
848 fin_rx.recv().expect("to receive final value");
849 });
850 }
851
852 async fn slow_pong(
853 fin_tx: std::sync::mpsc::Sender<()>,
854 mut tx: async_mpsc::Sender<u8>,
855 mut rx: async_mpsc::Receiver<u8>,
856 ) {
857 use zx::MonotonicDuration;
858 println!("starting pong!");
859 while let Some(next) = rx.next().await {
860 println!("pong! {next}");
861 fuchsia_async::Timer::new(fuchsia_async::MonotonicInstant::after(
862 MonotonicDuration::from_seconds(1),
863 ))
864 .await;
865 if next > 10 {
866 println!("bye!");
867 break;
868 }
869 tx.send(next + 1).await.unwrap();
870 }
871 fin_tx.send(()).unwrap();
872 }
873
874 #[test]
875 fn mixed_executor_async_ping_pong() {
876 with_raw_dispatcher("async ping pong", |dispatcher| {
877 let (fin_tx, fin_rx) = mpsc::channel();
878 let (ping_tx, pong_rx) = async_mpsc::channel(10);
879 let (pong_tx, ping_rx) = async_mpsc::channel(10);
880
881 dispatcher.spawn(ping(ping_tx, ping_rx));
883
884 let mut executor = fuchsia_async::LocalExecutor::default();
886 executor.run_singlethreaded(slow_pong(fin_tx, pong_tx, pong_rx));
887
888 fin_rx.recv().expect("to receive final value");
889 });
890 }
891
892 #[test]
893 fn wake_vector_registration() {
894 with_raw_dispatcher("wake vector test", |dispatcher| {
895 let dispatcher_ref =
896 DriverDispatcherRef::from_async_dispatcher(dispatcher.as_async_dispatcher_ref());
897 let event = zx::Event::create();
898 let signals = zx::Signals::USER_0;
899
900 let reg = dispatcher_ref.register_wake_vector(&event, signals).unwrap();
902 assert_eq!(reg.unregister(), Ok(()));
903
904 let reg_all =
907 dispatcher_ref.register_wake_vector(&event, zx::Signals::empty()).unwrap();
908 let reg_user0 = dispatcher_ref.register_wake_vector(&event, signals).unwrap();
909 assert_eq!(reg_all.unregister(), Ok(()));
910 assert_eq!(reg_user0.unregister(), Err(Status::NOT_FOUND));
911
912 {
914 let _reg_drop = dispatcher.register_wake_vector(&event, signals).unwrap();
915 }
916 });
917 }
918}