hyperactor/testing/
pingpong.rs1use std::time::Duration;
10
11use async_trait::async_trait;
12use hyperactor_config::Flattrs;
13use serde::Deserialize;
14use serde::Serialize;
15
16use crate as hyperactor; use crate::Actor;
18use crate::ActorRef;
19use crate::Context;
20use crate::Handler;
21use crate::Instance;
22use crate::OncePortRef;
23use crate::PortRef;
24use crate::RemoteSpawn;
25use crate::endpoint::Endpoint as _;
26use crate::mailbox::MessageEnvelope;
27use crate::mailbox::Undeliverable;
28use crate::mailbox::UndeliverableReason;
29
30#[derive(Serialize, Deserialize, Debug, typeuri::Named)]
35pub struct PingPongMessage(pub u64, pub ActorRef<PingPongActor>, pub OncePortRef<bool>);
36wirevalue::register_type!(PingPongMessage);
37
38#[derive(Debug)]
40#[hyperactor::export(PingPongMessage)]
41#[hyperactor::spawnable]
42pub struct PingPongActor {
43 undeliverable_port_ref: Option<PortRef<Undeliverable<MessageEnvelope>>>,
45 error_ttl: Option<u64>,
47 delay: Option<Duration>,
49}
50
51impl PingPongActor {
52 pub fn new(
58 undeliverable_port_ref: Option<PortRef<Undeliverable<MessageEnvelope>>>,
59 error_ttl: Option<u64>,
60 delay: Option<Duration>,
61 ) -> Self {
62 Self {
63 undeliverable_port_ref,
64 error_ttl,
65 delay,
66 }
67 }
68}
69
70#[async_trait]
71impl RemoteSpawn for PingPongActor {
72 type Params = (
73 Option<PortRef<Undeliverable<MessageEnvelope>>>,
74 Option<u64>,
75 Option<Duration>,
76 );
77
78 async fn new(
79 (undeliverable_port_ref, error_ttl, delay): Self::Params,
80 _environment: Flattrs,
81 ) -> anyhow::Result<Self> {
82 Ok(Self::new(undeliverable_port_ref, error_ttl, delay))
83 }
84}
85
86#[async_trait]
87impl Actor for PingPongActor {
88 async fn handle_delivery_failure_event(
89 &mut self,
90 cx: &Instance<Self>,
91 undelivered: Undeliverable<MessageEnvelope>,
92 ) -> Result<(), anyhow::Error> {
93 match &self.undeliverable_port_ref {
94 Some(port) => port.post(cx, undelivered),
95 None => crate::actor::handle_delivery_failure_event(self, cx, undelivered).await?,
96 }
97
98 Ok(())
99 }
100
101 async fn handle_undeliverable_message(
105 &mut self,
106 cx: &Instance<Self>,
107 _reason: UndeliverableReason,
108 undelivered: crate::mailbox::Undeliverable<crate::mailbox::MessageEnvelope>,
109 ) -> Result<(), anyhow::Error> {
110 match &self.undeliverable_port_ref {
111 Some(port) => {
112 port.post(cx, undelivered);
113 Ok(())
114 }
115 None => crate::actor::handle_undeliverable_message(cx, _reason, undelivered),
116 }
117 }
118
119 async fn handle_invalid_reference(
122 &mut self,
123 cx: &Instance<Self>,
124 invalid: crate::mailbox::InvalidReference,
125 undelivered: crate::mailbox::Undeliverable<crate::mailbox::MessageEnvelope>,
126 ) -> Result<(), anyhow::Error> {
127 match &self.undeliverable_port_ref {
128 Some(port) => {
129 port.post(cx, undelivered);
130 Ok(())
131 }
132 None => crate::actor::handle_invalid_reference(cx, invalid, undelivered),
133 }
134 }
135}
136
137#[async_trait]
138impl Handler<PingPongMessage> for PingPongActor {
139 async fn handle(
143 &mut self,
144 cx: &Context<Self>,
145 PingPongMessage(ttl, pong_actor, done_port): PingPongMessage,
146 ) -> anyhow::Result<()> {
147 if Some(ttl) == self.error_ttl {
150 anyhow::bail!("PingPong handler encountered an Error");
151 }
152 if ttl == 0 {
153 done_port.post(cx, true);
154 } else {
155 if let Some(delay) = self.delay {
156 tokio::time::sleep(delay).await;
157 }
158 let next_message = PingPongMessage(ttl - 1, cx.bind(), done_port);
159 pong_actor.post(cx, next_message);
160 }
161 Ok(())
162 }
163}