1use crate::actor::event_bus::{Answering, Event, GlobalEventBus};
2use crate::actor::traits::Handler;
3use crate::actor::{Cx, invoke_on_ui};
4use crate::trace;
5use once_cell::sync::Lazy;
6use parking_lot::RwLock;
7use std::any::{Any, TypeId};
8use std::collections::HashMap;
9use std::future::Future;
10use std::time::Duration;
11use tokio::sync::oneshot;
12use uuid::Uuid;
13
14tokio::task_local! {
15 static RPC_CHAIN: Vec<TypeId>;
23}
24
25fn current_rpc_chain() -> Vec<TypeId> {
26 RPC_CHAIN.try_with(|c| c.clone()).unwrap_or_default()
27}
28
29#[derive(Clone)]
30pub struct RpcRequest<T> {
31 pub correlation_id: Uuid,
32 pub payload: T,
33 pub chain: Vec<TypeId>,
38}
39
40impl<T: RpcCall> Event for RpcRequest<T> {
41 fn overheard(self) -> Self {
44 Self {
45 correlation_id: Uuid::nil(),
46 ..self
47 }
48 }
49}
50
51struct Unanswered(String);
53
54#[derive(Clone)]
55pub struct RpcResponse<T> {
56 pub correlation_id: Uuid,
57 pub payload: T,
58}
59
60impl<T: Clone + Send + 'static> Event for RpcResponse<T> {}
61
62pub trait RpcCall: Clone + Send + 'static {
70 type Response: Clone + Send + 'static;
71}
72
73impl<Req: RpcCall> RpcRequest<Req> {
74 pub fn reply(self, response: Req::Response) {
75 AsyncBus::reply(self.correlation_id, response);
76 }
77}
78
79pub struct Reply<T>(Answer<T>);
98
99enum Answer<T> {
100 Now(T),
101 Later(std::pin::Pin<Box<dyn Future<Output = T> + Send>>),
102}
103
104impl<T> Reply<T> {
105 pub fn now(value: T) -> Self {
106 Self(Answer::Now(value))
107 }
108
109 pub fn later(work: impl Future<Output = T> + Send + 'static) -> Self {
112 Self(Answer::Later(Box::pin(work)))
113 }
114}
115
116impl<T> From<T> for Reply<T> {
117 fn from(value: T) -> Self {
118 Self::now(value)
119 }
120}
121
122impl<T: Clone + Send + 'static> Reply<T> {
123 fn send(self, correlation_id: Uuid, chain: Vec<TypeId>) {
124 match self.0 {
125 Answer::Now(value) => AsyncBus::reply(correlation_id, value),
126 Answer::Later(work) => AsyncBus::spawn_reply(correlation_id, chain, work),
127 }
128 }
129}
130
131pub trait RpcHandler<Req: RpcCall>: 'static {
139 const DECLARED: Option<crate::actor::shape::Declared> = None;
141
142 fn handle_rpc(&mut self, req: Req, cx: Cx<Self, Req>) -> Reply<Req::Response>
143 where
144 Self: Sized;
145}
146
147impl<A, Req> Handler<RpcRequest<Req>> for A
148where
149 A: RpcHandler<Req> + 'static,
150 Req: RpcCall,
151{
152 const DECLARED: Option<crate::actor::shape::Declared> = <A as RpcHandler<Req>>::DECLARED;
153
154 const ANSWERS: bool = true;
155
156 fn handle(&mut self, msg: RpcRequest<Req>, cx: Cx<Self, RpcRequest<Req>>) {
157 self.handle_rpc(msg.payload, cx.handling())
158 .send(msg.correlation_id, msg.chain);
159 }
160}
161
162static PENDING_REQUESTS: Lazy<RwLock<HashMap<Uuid, oneshot::Sender<Box<dyn Any + Send>>>>> =
163 Lazy::new(|| RwLock::new(HashMap::new()));
164
165pub struct AsyncBus;
166
167impl AsyncBus {
168 pub async fn request<Req>(payload: Req, timeout: Duration) -> anyhow::Result<Req::Response>
169 where
170 Req: RpcCall,
171 {
172 let req_type = TypeId::of::<Req>();
173 let chain = current_rpc_chain();
174
175 if chain.contains(&req_type) {
185 panic!(
186 "AsyncBus: RPC cycle detected requesting {} - it (or a request that led back \
187 to it) is already awaiting its own reply {} level(s) up this call chain. This \
188 can never resolve: each hop is waiting on the next, all the way back to itself.",
189 std::any::type_name::<Req>(),
190 chain.len(),
191 );
192 }
193
194 let mut next_chain = chain;
195 next_chain.push(req_type);
196
197 let correlation_id = Uuid::new_v4();
198 let envelope = RpcRequest {
199 correlation_id,
200 payload,
201 chain: next_chain,
202 };
203
204 let (tx, rx) = oneshot::channel();
205 PENDING_REQUESTS.write().insert(correlation_id, tx);
206
207 let cause = trace::current();
214 invoke_on_ui(move || {
215 let _resumed = trace::resume(cause);
216 let bus = GlobalEventBus::instance();
217
218 let unanswered = match bus.answering::<RpcRequest<Req>>() {
219 Answering::Awake => return bus.publish(envelope),
220 Answering::Nobody => format!("nobody answers {}", std::any::type_name::<Req>()),
221 Answering::Asleep(answerer) => format!(
222 "{answerer}, which answers {}, is asleep",
223 std::any::type_name::<Req>()
224 ),
225 };
226
227 if let Some(tx) = PENDING_REQUESTS.write().remove(&correlation_id) {
228 let _ = tx.send(Box::new(Unanswered(unanswered)));
229 }
230 });
231
232 match tokio::time::timeout(timeout, rx).await {
233 Ok(Ok(any_res)) => match any_res.downcast::<RpcResponse<Req::Response>>() {
234 Ok(res) => Ok(res.payload),
235 Err(other) => match other.downcast::<Unanswered>() {
236 Ok(unanswered) => Err(anyhow::anyhow!(unanswered.0)),
237 Err(_) => Err(anyhow::anyhow!("Type mismatch in async response")),
238 },
239 },
240 Ok(Err(_)) => Err(anyhow::anyhow!("Response channel closed")),
241 Err(_) => {
242 PENDING_REQUESTS.write().remove(&correlation_id);
243 Err(anyhow::anyhow!("RPC request timed out"))
244 }
245 }
246 }
247
248 pub fn reply<Res>(correlation_id: Uuid, payload: Res)
249 where
250 Res: Clone + Send + 'static,
251 {
252 if correlation_id.is_nil() {
253 tracing::warn!(
254 reply = std::any::type_name::<Res>(),
255 "a reply from something that only hears the request went nowhere"
256 );
257 return;
258 }
259
260 let envelope = RpcResponse {
261 correlation_id,
262 payload,
263 };
264
265 if let Some(tx) = PENDING_REQUESTS.write().remove(&correlation_id) {
266 let _ = tx.send(Box::new(envelope.clone()));
267 }
268
269 GlobalEventBus::publish(envelope);
272 }
273
274 pub fn spawn_reply<Res, Fut>(correlation_id: Uuid, chain: Vec<TypeId>, fut: Fut)
297 where
298 Res: Clone + Send + 'static,
299 Fut: Future<Output = Res> + Send + 'static,
300 {
301 #[cfg(feature = "test-utils")]
302 let counted = crate::actor::event_bus::Counted::new();
303
304 crate::executor::spawn(RPC_CHAIN.scope(chain, async move {
305 #[cfg(feature = "test-utils")]
306 let _counted = counted;
307
308 let response = fut.await;
309 AsyncBus::reply(correlation_id, response);
310 }));
311 }
312}
313
314#[cfg(all(test, feature = "test-utils"))]
315mod tests {
316 use super::*;
317 use crate::actor::event_bus::EventBus;
318 use std::rc::Rc;
319 use std::sync::Mutex;
320 use std::sync::mpsc as std_mpsc;
321 use std::time::Duration as StdDuration;
322
323 static TEST_LOCK: Mutex<()> = Mutex::new(());
334
335 #[tokio::test]
336 async fn request_reply_round_trip_same_thread() {
337 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
338 #[derive(Clone, Debug, guinea_macros::Request)]
339 #[request(reply = Pong)]
340 struct Ping;
341 #[derive(Clone, Debug)]
342 struct Pong;
343
344 let _sub = GlobalEventBus::answer_fn(|_: Ping| Pong);
345
346 let handle = tokio::spawn(AsyncBus::request::<Ping>(Ping, StdDuration::from_secs(1)));
347 tokio::task::yield_now().await;
351 EventBus::process_queue();
352 EventBus::process_queue();
355
356 assert!(handle.await.unwrap().is_ok());
357 }
358
359 #[tokio::test]
360 async fn rpc_handler_replies_exactly_once_via_the_blanket_impl() {
361 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
362 use crate::actor::{Addr, UiThreadToken};
363
364 #[derive(Clone, Debug, guinea_macros::Request)]
365 #[request(reply = Echoed)]
366 struct Echo(u32);
367 #[derive(Clone, Debug)]
368 struct Echoed(u32);
369
370 struct EchoActor;
371 impl RpcHandler<Echo> for EchoActor {
372 fn handle_rpc(&mut self, Echo(n): Echo, _cx: Cx<Self, Echo>) -> Reply<Echoed> {
373 Reply::now(Echoed(n * 2))
374 }
375 }
376
377 let addr = Addr::new(EchoActor, UiThreadToken::dangerously_create_token_unchecked());
378 let _sub = GlobalEventBus::instance().subscribe::<EchoActor, RpcRequest<Echo>>(addr);
379
380 let handle = tokio::spawn(AsyncBus::request::<Echo>(Echo(21), StdDuration::from_secs(1)));
381 tokio::task::yield_now().await;
382 EventBus::process_queue();
386 EventBus::process_queue();
387
388 let response = handle.await.unwrap().expect("rpc handler should have replied");
389 assert_eq!(response.0, 42);
390 }
391
392 #[tokio::test]
393 async fn a_request_nobody_answers_fails_at_once() {
394 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
395 #[derive(Clone, Debug, guinea_macros::Request)]
396 #[request(reply = NeverReplied)]
397 struct Unasked;
398 #[derive(Clone, Debug)]
399 struct NeverReplied;
400
401 let _hears = GlobalEventBus::instance().subscribe_fn(|_: RpcRequest<Unasked>| {});
402
403 let handle = tokio::spawn(AsyncBus::request::<Unasked>(Unasked, StdDuration::from_secs(60)));
404 tokio::task::yield_now().await;
405 EventBus::process_queue();
406
407 let error = tokio::time::timeout(StdDuration::from_secs(5), handle)
408 .await
409 .expect("it waited for a reply nothing could give")
410 .unwrap()
411 .expect_err("nothing answers it");
412 assert!(error.to_string().contains("nobody answers"), "{error}");
413 }
414
415 #[test]
416 #[should_panic(expected = "a request has exactly one answerer")]
417 fn a_second_answerer_is_refused_when_it_subscribes() {
418 #[derive(Clone, Debug, guinea_macros::Request)]
419 #[request(reply = Answer)]
420 struct Question;
421 #[derive(Clone, Debug)]
422 struct Answer;
423
424 let bus = Rc::new(EventBus::new());
425 let _first = bus.answer_fn(|_: Question| Answer);
426 let _second = bus.answer_fn(|_: Question| Answer);
427 }
428
429 #[test]
430 #[should_panic(expected = "a request has exactly one answerer")]
431 fn two_actors_returning_the_reply_cannot_both_answer() {
432 use crate::actor::{Addr, UiThreadToken};
433
434 #[derive(Clone, Debug, guinea_macros::Request)]
435 #[request(reply = Answer)]
436 struct Question;
437 #[derive(Clone, Debug)]
438 struct Answer;
439
440 struct Service;
441 struct Monitor;
442 impl RpcHandler<Question> for Service {
443 fn handle_rpc(&mut self, _: Question, _cx: Cx<Self, Question>) -> Reply<Answer> {
444 Reply::now(Answer)
445 }
446 }
447 impl RpcHandler<Question> for Monitor {
448 fn handle_rpc(&mut self, _: Question, _cx: Cx<Self, Question>) -> Reply<Answer> {
449 Reply::now(Answer)
450 }
451 }
452
453 let token = UiThreadToken::dangerously_create_token_unchecked();
454 let bus = Rc::new(EventBus::new());
455 let _service =
456 bus.subscribe::<Service, RpcRequest<Question>>(Addr::new(Service, token.clone()));
457 let _monitor =
458 bus.subscribe::<Monitor, RpcRequest<Question>>(Addr::new(Monitor, token));
459 }
460
461 #[test]
462 fn the_answerer_s_place_is_free_again_once_it_goes() {
463 #[derive(Clone, Debug, guinea_macros::Request)]
464 #[request(reply = Answer)]
465 struct Question;
466 #[derive(Clone, Debug)]
467 struct Answer;
468
469 let bus = Rc::new(EventBus::new());
470 drop(bus.answer_fn(|_: Question| Answer));
471 let _next = bus.answer_fn(|_: Question| Answer);
472 }
473
474 #[tokio::test]
475 async fn a_listener_hears_the_request_and_its_reply_goes_nowhere() {
476 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
477 use std::sync::Arc;
478 use std::sync::atomic::{AtomicBool, Ordering};
479
480 #[derive(Clone, Debug, guinea_macros::Request)]
481 #[request(reply = Count)]
482 struct HowMany;
483 #[derive(Clone, Debug, PartialEq)]
484 struct Count(u32);
485
486 let heard = Arc::new(AtomicBool::new(false));
487 let hearing = heard.clone();
488 let _listener = GlobalEventBus::subscribe_fn(move |request: RpcRequest<HowMany>| {
489 hearing.store(true, Ordering::SeqCst);
490 request.reply(Count(0));
491 });
492 let _answerer = GlobalEventBus::answer_fn(|_: HowMany| Count(7));
493
494 let handle = tokio::spawn(AsyncBus::request::<HowMany>(HowMany, StdDuration::from_secs(1)));
495 tokio::task::yield_now().await;
496 EventBus::process_queue();
497 EventBus::process_queue();
498
499 assert_eq!(handle.await.unwrap().unwrap(), Count(7));
500 assert!(heard.load(Ordering::SeqCst), "the listener was not told");
501 }
502
503 #[tokio::test]
504 async fn a_request_whose_answerer_sleeps_fails_at_once() {
505 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
506 use crate::actor::{Addr, UiThreadToken};
507 use crate::scope::ScopeTree;
508
509 #[derive(Clone, Debug, guinea_macros::Request)]
510 #[request(reply = Done)]
511 struct Work;
512 #[derive(Clone, Debug)]
513 struct Done;
514
515 struct Worker;
516 impl RpcHandler<Work> for Worker {
517 fn handle_rpc(&mut self, _: Work, _cx: Cx<Self, Work>) -> Reply<Done> {
518 Reply::now(Done)
519 }
520 }
521
522 let scope = ScopeTree::new();
523 let addr = Addr::new(Worker, UiThreadToken::dangerously_create_token_unchecked());
524 addr.live_in(scope.scope(), Some(&Rc::new(EventBus::new())));
525 let _sub = GlobalEventBus::instance().subscribe::<Worker, RpcRequest<Work>>(addr);
526 scope.sleep();
527
528 let handle = tokio::spawn(AsyncBus::request::<Work>(Work, StdDuration::from_secs(60)));
529 tokio::task::yield_now().await;
530 EventBus::process_queue();
531
532 let error = tokio::time::timeout(StdDuration::from_secs(5), handle)
533 .await
534 .expect("it waited for an answerer that sleeps")
535 .unwrap()
536 .expect_err("its answerer sleeps");
537 assert!(error.to_string().contains("is asleep"), "{error}");
538 }
539
540 #[test]
551 fn request_resolves_when_subscriber_is_on_a_different_os_thread() {
552 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
553 #[derive(Clone, Debug, guinea_macros::Request)]
554 #[request(reply = Pong)]
555 struct Ping;
556 #[derive(Clone, Debug)]
557 struct Pong;
558
559 let (ready_tx, ready_rx) = std_mpsc::channel::<()>();
560 let (stop_tx, stop_rx) = std_mpsc::channel::<()>();
561
562 let ui_thread = std::thread::spawn(move || {
567 let _sub = GlobalEventBus::answer_fn(|_: Ping| Pong);
568 ready_tx.send(()).unwrap();
569 while stop_rx.try_recv().is_err() {
570 EventBus::process_queue();
571 std::thread::sleep(StdDuration::from_millis(5));
572 }
573 });
574
575 ready_rx.recv().unwrap();
576
577 let rt = tokio::runtime::Builder::new_current_thread()
579 .enable_time()
580 .build()
581 .unwrap();
582 let result =
583 rt.block_on(AsyncBus::request::<Ping>(Ping, StdDuration::from_secs(2)));
584
585 stop_tx.send(()).unwrap();
586 ui_thread.join().unwrap();
587
588 assert!(
589 result.is_ok(),
590 "expected a reply delivered from another OS thread, got {result:?}"
591 );
592 }
593
594 #[tokio::test]
598 async fn handler_macro_rpc_heuristic_sync_and_async() {
599 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
600 use crate::actor::{Addr, AsyncContext, UiThreadToken};
601 use guinea_macros::handler;
602
603 #[derive(Clone, Debug, guinea_macros::Request)]
604 #[request(reply = Doubled)]
605 struct Double(u32);
606 #[derive(Clone, Debug)]
607 struct Doubled(u32);
608
609 #[derive(Clone, Debug, guinea_macros::Request)]
610 #[request(reply = Sum)]
611 struct DelayedAdd(u32, u32);
612 #[derive(Clone, Debug)]
613 struct Sum(u32);
614
615 struct MathActor;
616
617 #[handler]
619 fn double(this: &mut MathActor, Double(n): Double) -> Doubled {
620 let _ = this;
621 Doubled(n * 2)
622 }
623
624 #[handler]
627 async fn delayed_add(_ctx: AsyncContext<MathActor>, req: DelayedAdd) -> Sum {
628 tokio::time::sleep(StdDuration::from_millis(1)).await;
629 Sum(req.0 + req.1)
630 }
631
632 let addr =
633 Addr::new(MathActor, UiThreadToken::dangerously_create_token_unchecked());
634 let _sub_double =
635 GlobalEventBus::instance().subscribe::<MathActor, RpcRequest<Double>>(addr.clone());
636 let _sub_add = GlobalEventBus::instance()
637 .subscribe::<MathActor, RpcRequest<DelayedAdd>>(addr.clone());
638
639 let double_handle =
640 tokio::spawn(AsyncBus::request::<Double>(Double(21), StdDuration::from_secs(1)));
641 tokio::task::yield_now().await;
642 EventBus::process_queue(); EventBus::process_queue(); assert_eq!(double_handle.await.unwrap().unwrap().0, 42);
645
646 let add_handle = tokio::spawn(AsyncBus::request::<DelayedAdd>(
647 DelayedAdd(2, 3),
648 StdDuration::from_secs(1),
649 ));
650 tokio::task::yield_now().await;
651 EventBus::process_queue(); for _ in 0..20 {
656 tokio::time::sleep(StdDuration::from_millis(5)).await;
657 EventBus::process_queue();
658 }
659 assert_eq!(add_handle.await.unwrap().unwrap().0, 5);
660 }
661
662 #[tokio::test]
681 async fn cross_actor_rpc_cycle_is_detected_immediately() {
682 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
683 use crate::actor::{Addr, AsyncContext, UiThreadToken};
684 use guinea_macros::handler;
685 use std::panic;
686 use std::sync::{Arc, Mutex};
687
688 #[derive(Clone, Debug, guinea_macros::Request)]
689 #[request(reply = RespA)]
690 struct ReqA(u32);
691 #[derive(Clone, Debug)]
692 struct RespA(u32);
693
694 #[derive(Clone, Debug, guinea_macros::Request)]
695 #[request(reply = RespB)]
696 struct ReqB(u32);
697 #[derive(Clone, Debug)]
698 struct RespB(u32);
699
700 struct ActorA;
701 struct ActorB;
702
703 #[handler]
704 async fn handle_req_a(_ctx: AsyncContext<ActorA>, req: ReqA) -> RespA {
705 let RespB(n) = AsyncBus::request::<ReqB>(ReqB(req.0), StdDuration::from_secs(5))
706 .await
707 .expect("should never resolve normally - the cycle panics first");
708 RespA(n)
709 }
710
711 #[handler]
712 async fn handle_req_b(_ctx: AsyncContext<ActorB>, req: ReqB) -> RespB {
713 let RespA(n) = AsyncBus::request::<ReqA>(ReqA(req.0), StdDuration::from_secs(5))
715 .await
716 .expect("should never resolve normally - the cycle panics first");
717 RespB(n)
718 }
719
720 let panics: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(Vec::new()));
725 let hook_slot = panics.clone();
726 let prev_hook = panic::take_hook();
727 panic::set_hook(Box::new(move |info| {
728 hook_slot.lock().unwrap().push(info.to_string());
729 }));
730
731 let seen_the_cycle = || {
732 panics
733 .lock()
734 .unwrap()
735 .iter()
736 .any(|message| message.contains("RPC cycle detected"))
737 };
738
739 let addr_a = Addr::new(ActorA, UiThreadToken::dangerously_create_token_unchecked());
740 let addr_b = Addr::new(ActorB, UiThreadToken::dangerously_create_token_unchecked());
741 let _sub_a = GlobalEventBus::instance().subscribe::<ActorA, RpcRequest<ReqA>>(addr_a);
742 let _sub_b = GlobalEventBus::instance().subscribe::<ActorB, RpcRequest<ReqB>>(addr_b);
743
744 let handle = tokio::spawn(AsyncBus::request::<ReqA>(ReqA(1), StdDuration::from_secs(5)));
745
746 let deadline = std::time::Instant::now() + StdDuration::from_secs(5);
752 while std::time::Instant::now() < deadline {
753 tokio::task::yield_now().await;
754 EventBus::process_queue();
755 if seen_the_cycle() {
756 break;
757 }
758 tokio::time::sleep(StdDuration::from_millis(1)).await;
759 }
760
761 panic::set_hook(prev_hook);
762 handle.abort(); assert!(
765 seen_the_cycle(),
766 "expected AsyncBus's cycle check to panic inside actor B's handler; \
767 the panics seen were: {:?}",
768 panics.lock().unwrap()
769 );
770 }
771
772 #[test]
785 fn many_concurrent_requests_do_not_deadlock() {
786 let _guard = TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner());
787 use crate::actor::{Addr, UiThreadToken};
788 use std::sync::Arc;
789 use std::sync::atomic::{AtomicBool, Ordering};
790
791 #[derive(Clone, Debug, guinea_macros::Request)]
792 #[request(reply = Sum)]
793 struct Add(u32, u32);
794 #[derive(Clone, Debug)]
795 struct Sum(u32);
796
797 struct AddActor;
798 impl RpcHandler<Add> for AddActor {
799 fn handle_rpc(&mut self, Add(a, b): Add, _cx: Cx<Self, Add>) -> Reply<Sum> {
800 Reply::now(Sum(a + b))
801 }
802 }
803
804 let done = Arc::new(AtomicBool::new(false));
805 let watchdog_done = done.clone();
806 let watchdog = std::thread::spawn(move || {
807 for _ in 0..100 {
808 if watchdog_done.load(Ordering::SeqCst) {
809 return;
810 }
811 std::thread::sleep(StdDuration::from_millis(50));
812 }
813 eprintln!(
814 "many_concurrent_requests_do_not_deadlock: deadline exceeded \
815 without completing - treating this as a deadlock and \
816 aborting instead of hanging"
817 );
818 std::process::abort();
819 });
820
821 let (ready_tx, ready_rx) = std_mpsc::channel::<()>();
822 let (stop_tx, stop_rx) = std_mpsc::channel::<()>();
823
824 let ui_thread = std::thread::spawn(move || {
828 let addr =
829 Addr::new(AddActor, UiThreadToken::dangerously_create_token_unchecked());
830 let _sub = GlobalEventBus::instance().subscribe::<AddActor, RpcRequest<Add>>(addr);
831 ready_tx.send(()).unwrap();
832 while stop_rx.try_recv().is_err() {
833 EventBus::process_queue();
834 std::thread::sleep(StdDuration::from_millis(2));
835 }
836 });
837
838 ready_rx.recv().unwrap();
839
840 const REQUESTER_THREADS: u32 = 8;
841 const REQUESTS_PER_THREAD: u32 = 25;
842
843 let requesters: Vec<_> = (0..REQUESTER_THREADS)
844 .map(|t| {
845 std::thread::spawn(move || {
846 let rt = tokio::runtime::Builder::new_current_thread()
847 .enable_time()
848 .build()
849 .unwrap();
850 rt.block_on(async {
851 for i in 0..REQUESTS_PER_THREAD {
852 let result =
853 AsyncBus::request::<Add>(Add(t, i), StdDuration::from_secs(5))
854 .await
855 .unwrap_or_else(|e| panic!("request {t}/{i} failed: {e}"));
856 assert_eq!(result.0, t + i);
857 }
858 });
859 })
860 })
861 .collect();
862
863 for r in requesters {
864 r.join().unwrap();
865 }
866
867 stop_tx.send(()).unwrap();
868 ui_thread.join().unwrap();
869
870 done.store(true, Ordering::SeqCst);
873 watchdog.join().unwrap();
874 }
875}
876
877#[deprecated(
878 since = "0.18.6",
879 note = "write the reply on the request: `#[derive(guinea::Request)] #[request(reply = Res)]`"
880)]
881#[macro_export]
882macro_rules! rpc_bind {
883 ($( $req:ident => $res:ident );* $(;)?) => {
884 $(
885 impl $crate::actor::event_bus::rpc::RpcCall for $req {
886 type Response = $res;
887 }
888 )*
889 };
890}