use super::*;
struct FakeNodeSession;
impl FakeNodeSession {
async fn get_node_session() -> (ActorRef<crate::node::NodeSessionMessage>, JoinHandle<()>) {
let (r, h) = Actor::spawn(None, FakeNodeSession, ())
.await
.expect("Failed to start fake node session");
let cell: ractor::ActorCell = r.into();
(cell.into(), h)
}
}
#[cfg_attr(feature = "async-trait", ractor::async_trait)]
impl Actor for FakeNodeSession {
type Msg = crate::node::NodeSessionMessage;
type State = ();
type Arguments = ();
async fn pre_start(
&self,
_: ActorRef<Self::Msg>,
_: Self::Arguments,
) -> Result<Self::State, ActorProcessingErr> {
Ok(())
}
}
#[ractor::concurrency::test]
async fn remote_actor_serialized_message_handling() {
let (actor, handle) = FakeNodeSession::get_node_session().await;
let (remote_actor_ref, remote_actor_handle) = Actor::spawn(None, RemoteActor, actor.clone())
.await
.expect("Failed to spawn remote actor");
let remote_actor_instance = RemoteActor;
let mut remote_actor_state = RemoteActorState::new(actor.clone());
let bad_handler = remote_actor_instance
.handle(
remote_actor_ref.clone(),
RemoteActorMessage,
&mut remote_actor_state,
)
.await;
assert!(bad_handler.is_err());
let cast = SerializedMessage::Cast {
variant: "A".to_string(),
args: vec![1, 2, 3],
metadata: None,
};
let cast_output = remote_actor_instance
.handle_serialized(remote_actor_ref.clone(), cast, &mut remote_actor_state)
.await;
assert!(cast_output.is_ok());
assert_eq!(0, remote_actor_state.message_tag);
let (tx, _rx) = ractor::concurrency::oneshot();
let call = SerializedMessage::Call {
variant: "B".to_string(),
args: vec![1, 2, 3],
reply: tx.into(),
metadata: None,
};
let call_output = remote_actor_instance
.handle_serialized(remote_actor_ref.clone(), call, &mut remote_actor_state)
.await;
assert!(call_output.is_ok());
assert_eq!(1, remote_actor_state.message_tag);
assert!(remote_actor_state.pending_requests.contains_key(&1));
let reply = SerializedMessage::CallReply(1, vec![3, 4, 5]);
let reply_output = remote_actor_instance
.handle_serialized(remote_actor_ref.clone(), reply, &mut remote_actor_state)
.await;
assert!(reply_output.is_ok());
assert!(!remote_actor_state.pending_requests.contains_key(&1));
remote_actor_ref.stop(None);
actor.stop(None);
remote_actor_handle.await.unwrap();
handle.await.unwrap();
}
#[ractor::concurrency::test]
async fn abandoned_requests_are_reclaimed_incrementally() {
let (session, session_handle) = FakeNodeSession::get_node_session().await;
let (remote_actor, remote_actor_handle) = Actor::spawn(None, RemoteActor, session.clone())
.await
.unwrap();
let mut state = RemoteActorState::new(session.clone());
let (reply, receiver) = ractor::concurrency::oneshot::<Vec<u8>>();
drop(receiver);
RemoteActor
.handle_serialized(
remote_actor.clone(),
SerializedMessage::Call {
variant: "Call".to_string(),
args: vec![],
reply: reply.into(),
metadata: None,
},
&mut state,
)
.await
.unwrap();
assert_eq!(state.pending_requests.len(), 1);
RemoteActor
.handle_serialized(
remote_actor.clone(),
SerializedMessage::Cast {
variant: "Cast".to_string(),
args: vec![],
metadata: None,
},
&mut state,
)
.await
.unwrap();
assert!(state.pending_requests.is_empty());
remote_actor.stop(None);
session.stop(None);
remote_actor_handle.await.unwrap();
session_handle.await.unwrap();
}
#[ractor::concurrency::test]
async fn forwarding_failure_releases_the_reply_port() {
let (session, session_handle) = FakeNodeSession::get_node_session().await;
session.stop(None);
session_handle.await.unwrap();
let (remote_actor, remote_actor_handle) = Actor::spawn(None, RemoteActor, session.clone())
.await
.unwrap();
let mut state = RemoteActorState::new(session);
let (reply, receiver) = ractor::concurrency::oneshot::<Vec<u8>>();
RemoteActor
.handle_serialized(
remote_actor.clone(),
SerializedMessage::Call {
variant: "Call".to_string(),
args: vec![],
reply: reply.into(),
metadata: None,
},
&mut state,
)
.await
.unwrap();
assert!(state.pending_requests.is_empty());
assert!(receiver.await.is_err());
remote_actor.stop(None);
remote_actor_handle.await.unwrap();
}