use crate::environment::FetchEnvironmentEvent;
use crate::{
capabilities::SendCapabilityRequest, inventory::RefreshInventoryEvent,
transport::http_handler::login_to_simulator,
};
use actix::prelude::*;
use actix_rt::time;
use benthic_protocol::errors::SessionError;
use benthic_protocol::messages::ui::errors::MailboxSessionError;
#[cfg(feature = "environment")]
use benthic_protocol::session::EnvironmentCache;
#[cfg(feature = "inventory")]
use benthic_protocol::session::InventoryData;
use benthic_protocol::session::RegionData;
use benthic_protocol::{
messages::ui::{
errors::{CapabilityError, CircuitCodeError, CompleteAgentMovementError, FeatureError},
login_event::Login,
login_response::LoginResponse,
ui_messages::{UIMessage, UIResponse},
water_update::WaterUpdate,
},
session::{ServerState, Session},
};
use glam::Vec2;
use log::{error, info};
use metaverse_avatar::avatar::Avatar;
use metaverse_cache::initialize_sqlite::init_sqlite;
use metaverse_environment::land::Land;
use metaverse_messages::{
http::capabilities::{Capability, CapabilityRequest},
packet::packet_protocol::Packet,
udp::{
agent::agent_update::AgentUpdate,
chat::chat_from_viewer::ChatFromViewer,
core::{
agent_throttle::{AgentThrottle, ThrottleData},
circuit_code::CircuitCode,
complete_agent_movement::CompleteAgentMovementData,
complete_ping_check::CompletePingCheck,
logout_request::LogoutRequest,
packet_ack::PacketAck,
region_handshake::RegionHandshake,
region_handshake_reply::RegionHandshakeReply,
},
},
};
use rgb::Rgba;
use std::path::Path;
use std::{
collections::{HashMap, HashSet},
net::UdpSocket as SyncUdpSocket,
path::PathBuf,
sync::{Arc, Mutex},
thread::sleep,
};
use tokio::{net::UdpSocket, sync::Notify, time::Duration};
#[derive(Debug)]
pub struct Mailbox {
pub client_socket: u16,
pub server_to_ui_socket: String,
pub inventory_db_location: PathBuf,
pub server_acks: HashSet<u32>,
pub viewer_acks: HashSet<u32>,
pub state: Arc<Mutex<ServerState>>,
pub notify: Arc<Notify>,
pub session: Option<Session<Capability, Avatar, Land, UdpSocket>>,
pub sent_packet_count: u16,
pub ping_info: PingInfo,
}
impl Mailbox {
pub fn set_state(&mut self, new_state: ServerState, _ctx: &mut Context<Self>) {
let state_clone: Arc<Mutex<ServerState>> = Arc::clone(&self.state);
{
let mut state = state_clone.lock().unwrap();
*state = new_state.clone();
}
if new_state == ServerState::Running || new_state == ServerState::Stopped {
self.notify.notify_one();
}
}
}
impl Actor for Mailbox {
type Context = Context<Self>;
fn started(&mut self, ctx: &mut Self::Context) {
info!("Actix Mailbox has started");
self.set_state(ServerState::Running, ctx);
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct StartSession {
pub session: Session<Capability, Avatar, Land, UdpSocket>,
}
impl Handler<StartSession> for Mailbox {
type Result = ();
fn handle(&mut self, msg: StartSession, ctx: &mut Self::Context) -> Self::Result {
let mut session = msg.session;
if let Some(current_session) = self.session.as_ref() {
session.socket = current_session.socket.clone();
}
self.session = Some(session);
if let Some(session) = self.session.as_ref()
&& session.socket.is_none()
{
let addr = format!("{}:{}", session.local_ip, self.client_socket);
let addr_clone = addr.clone();
let mailbox_addr = ctx.address();
info!("session established, starting UDP processing");
let fut = async move {
match UdpSocket::bind(&addr).await {
Ok(sock) => {
info!("Successfully bound to {}", addr);
let sock = Arc::new(sock);
tokio::spawn(Mailbox::start_udp_read(sock.clone(), mailbox_addr));
Ok(sock) }
Err(e) => {
error!("Failed to bind to {}: {}", addr_clone, e);
Err(e)
}
}
};
ctx.spawn(fut.into_actor(self).map(|result, act, _| match result {
Ok(sock) => {
if let Some(session) = &mut act.session {
session.socket = Some(sock);
}
}
Err(_) => {
panic!("Socket binding failed");
}
}));
}
}
}
#[derive(Debug)]
pub struct PingInfo {
pub ping_number: u8,
pub ping_latency: Duration,
pub last_ping: time::Instant,
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct HandlePing {
pub ping_id: u8,
}
impl Handler<HandlePing> for Mailbox {
type Result = ();
fn handle(&mut self, msg: HandlePing, ctx: &mut Self::Context) -> Self::Result {
ctx.address().do_send(OutgoingPacket {
packet: Packet::new_complete_ping_check(CompletePingCheck {
ping_id: msg.ping_id,
}),
});
self.ping_info.ping_latency = time::Instant::now() - self.ping_info.last_ping;
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct SendAckList {}
impl Handler<SendAckList> for Mailbox {
type Result = ();
fn handle(&mut self, _: SendAckList, ctx: &mut Self::Context) -> Self::Result {
if let Some(ref session) = self.session {
if self.server_acks.is_empty() {
return;
}
let packet_ids: Vec<u32> = self.server_acks.drain().collect();
let addr = session.address.clone();
let packet = Packet::new_packet_ack(PacketAck { packet_ids }).to_bytes();
let sock_clone = session.socket.clone().unwrap();
let ack_wait = async move {
if let Err(e) = sock_clone.send_to(&packet, addr).await {
println!("Failed to send ack: {:?}", e)
};
};
ctx.spawn(ack_wait.into_actor(self));
}
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct AddToAckList {
pub id: u32,
}
impl Handler<AddToAckList> for Mailbox {
type Result = ();
fn handle(&mut self, msg: AddToAckList, ctx: &mut Self::Context) -> Self::Result {
self.server_acks.insert(msg.id);
ctx.address().do_send(SendAckList {});
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct ResendPacket {
pub packet: Packet,
}
impl Handler<ResendPacket> for Mailbox {
type Result = ();
fn handle(&mut self, mut msg: ResendPacket, ctx: &mut Self::Context) -> Self::Result {
if self
.viewer_acks
.contains(&msg.packet.header.sequence_number)
{
msg.packet.header.resent = true;
ctx.address().do_send(OutgoingPacket { packet: msg.packet });
}
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct OutgoingPacket {
pub packet: Packet,
}
impl Handler<OutgoingPacket> for Mailbox {
type Result = ();
fn handle(&mut self, mut msg: OutgoingPacket, ctx: &mut Self::Context) -> Self::Result {
if let Some(session) = self.session.as_mut() {
let addr = session.address.clone();
if !msg.packet.header.resent {
msg.packet.header.sequence_number = session.sequence_number as u32;
session.sequence_number += 1;
}
let data = msg.packet.to_bytes().clone();
let socket_clone = session.socket.as_ref().unwrap().clone();
let fut = async move {
if let Err(e) = socket_clone.send_to(&data, &addr).await {
error!("Failed to send data: {}", e);
}
};
ctx.spawn(fut.into_actor(self));
if msg.packet.header.reliable {
self.viewer_acks.insert(msg.packet.header.sequence_number);
ctx.notify_later(ResendPacket { packet: msg.packet }, Duration::from_secs(1));
};
}
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct HandleRegionHandshake {
pub region_handshake: RegionHandshake,
}
impl Handler<HandleRegionHandshake> for Mailbox {
type Result = ();
fn handle(&mut self, msg: HandleRegionHandshake, ctx: &mut Self::Context) -> Self::Result {
if let Some(session) = &mut self.session {
session.region_data.water_height = msg.region_handshake.water_height;
ctx.address().do_send(SendUIMessage {
ui_message: UIMessage::new_water_update(WaterUpdate {
height: msg.region_handshake.water_height,
color: Rgba::from((0.0f32, 94.0, 184.0, 0.5)),
}),
});
ctx.address().do_send(OutgoingPacket {
packet: Packet::new_region_handshake_reply(RegionHandshakeReply {
session_id: session.session_id,
agent_id: session.agent_id,
flags: 0,
}),
});
}
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct HandlePacketAck {
pub packet_ack: PacketAck,
}
impl Handler<HandlePacketAck> for Mailbox {
type Result = ();
fn handle(&mut self, msg: HandlePacketAck, _ctx: &mut Self::Context) -> Self::Result {
for id in msg.packet_ack.packet_ids {
self.viewer_acks.remove(&id);
}
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct HandleUIResponse {
pub ui_response: UIResponse,
}
impl Handler<HandleUIResponse> for Mailbox {
type Result = ();
fn handle(&mut self, msg: HandleUIResponse, ctx: &mut Self::Context) -> Self::Result {
if let UIResponse::Login(data) = &msg.ui_response {
let ctx_addr = ctx.address().clone();
let login_data = data.clone();
let db_location = self.inventory_db_location.clone();
actix::spawn(async move {
if let Err(e) = handle_login(login_data, &ctx_addr, &db_location).await {
error!("{:?}", e);
};
});
}
if let Some(ref session) = self.session {
match msg.ui_response {
UIResponse::ChatFromViewer(data) => {
ctx.address().do_send(OutgoingPacket {
packet: Packet::new_chat_from_viewer(ChatFromViewer {
session_id: session.session_id,
agent_id: session.agent_id,
channel: data.channel,
message: data.message,
message_type: data.message_type,
}),
});
}
UIResponse::Logout(_) => {
ctx.address().do_send(OutgoingPacket {
packet: Packet::new_logout_request(LogoutRequest {
session_id: session.session_id,
agent_id: session.agent_id,
}),
});
}
UIResponse::AgentUpdate(mut data) => {
data.agent_id = session.agent_id;
data.session_id = session.session_id;
ctx.address().do_send(OutgoingPacket {
packet: Packet::new_agent_update(AgentUpdate {
session_id: data.session_id,
agent_id: data.session_id,
body_rotation: data.body_rotation,
head_rotation: data.head_rotation,
state: data.state,
camera_center: data.camera_center,
camera_at_axis: data.camera_at_axis,
camera_left_axis: data.camera_left_axis,
camera_up_axis: data.camera_up_axis,
far: data.far,
control_flags: data.control_flags,
flags: data.flags,
}),
});
}
data => {
error!("Unrecognized UIMessage: {:?}", data)
}
}
}
}
}
#[derive(Debug, Message)]
#[rtype(result = "()")]
pub struct SendUIMessage {
pub ui_message: UIMessage,
}
impl Handler<SendUIMessage> for Mailbox {
type Result = ();
fn handle(&mut self, msg: SendUIMessage, _: &mut Self::Context) -> Self::Result {
let client_socket = SyncUdpSocket::bind("0.0.0.0:0").unwrap();
if let Err(e) = client_socket.send_to(&msg.ui_message.to_bytes(), &self.server_to_ui_socket)
{
error!("Failed to send UI message:{:?}", e)
}
}
}
async fn handle_login(
login_data: Login,
mailbox_addr: &actix::Addr<Mailbox>,
db_path: &Path,
) -> Result<(), SessionError> {
let (login_response, local_ip) = match login_to_simulator(login_data).await {
Ok((login_response, local_ip)) => {
if let Err(e) = mailbox_addr
.send(SendUIMessage {
ui_message: UIMessage::new_login_response_event(LoginResponse {
firstname: login_response.first_name.clone(),
lastname: login_response.last_name.clone(),
agent_id: login_response.agent_id,
}),
})
.await
{
error!("Failed to send login response to UI {:?}", e)
};
(login_response, local_ip)
}
Err(e) => {
if let Err(e) = mailbox_addr
.send(SendUIMessage {
ui_message: UIMessage::new_session_error(e.clone().into()),
})
.await
{
error!("Failed to send session error: {:?}", e)
};
Err(e)?
}
};
let connection = init_sqlite(db_path.to_path_buf())
.await
.map_err(|e| FeatureError::Inventory(format!("Failed to initialize SQLite: {}", e)))?;
if let Err(e) = mailbox_addr
.send(StartSession {
session: Session {
inventory_db_connection: connection,
agent_id: login_response.agent_id,
session_id: login_response.session_id,
address: format!("{}:{}", login_response.sim_ip, login_response.sim_port),
seed_capability_url: login_response.seed_capability.unwrap(),
sequence_number: 0,
local_ip,
capability_urls: HashMap::new(),
region_data: RegionData {
region_coordinates: Vec2 {
x: (login_response.region_x.unwrap() as f32),
y: (login_response.region_y.unwrap() as f32),
},
region_id: format!("{}:{}", login_response.sim_ip, login_response.sim_port),
..Default::default()
},
environment_cache: EnvironmentCache {
patch_queue: HashMap::new(),
patch_cache: HashMap::new(),
},
inventory_data: InventoryData {
inventory_root: login_response.inventory_root.ok_or_else(|| {
FeatureError::Inventory(
"Login response contained no inventory_root".to_string(),
)
})?,
inventory_lib_owner: login_response.inventory_lib_owner.ok_or_else(|| {
FeatureError::Inventory(
"Login response contained no inventory_lib_owner".to_string(),
)
})?,
inventory_init: false,
},
socket: None,
#[cfg(feature = "avatar")]
avatars: HashMap::new(),
},
})
.await
{
Err(MailboxSessionError {
message: e.to_string(),
})?;
};
match CapabilityRequest::new_capability_request(vec![
Capability::ViewerAsset,
Capability::FetchInventoryDescendents2,
Capability::ExtEnvironment,
]) {
Ok(caps) => {
if let Err(e) = mailbox_addr
.send(SendCapabilityRequest {
capability_request: caps,
})
.await
{
Err(CapabilityError {
message: e.to_string(),
})?
}
}
Err(e) => {
error!("{:?}", e)
}
};
if let Err(e) = mailbox_addr
.send(OutgoingPacket {
packet: Packet::new_circuit_code(CircuitCode {
code: login_response.circuit_code,
session_id: login_response.session_id,
id: login_response.agent_id,
}),
})
.await
{
Err(CircuitCodeError {
message: e.to_string(),
})?;
};
sleep(Duration::from_secs(1));
if let Err(e) = mailbox_addr
.send(OutgoingPacket {
packet: Packet::new_complete_agent_movement(CompleteAgentMovementData {
circuit_code: login_response.circuit_code,
session_id: login_response.session_id,
agent_id: login_response.agent_id,
}),
})
.await
{
Err(CompleteAgentMovementError {
message: e.to_string(),
})?;
};
if let Err(e) = mailbox_addr
.send(OutgoingPacket {
packet: Packet::new_agent_throttle(AgentThrottle {
agent_id: login_response.agent_id,
session_id: login_response.session_id,
circuit_code: login_response.circuit_code,
gen_counter: 0,
throttles: ThrottleData {
..Default::default()
},
}),
})
.await
{
Err(CompleteAgentMovementError {
message: e.to_string(),
})?;
};
#[cfg(feature = "inventory")]
if let Err(e) = mailbox_addr
.send(RefreshInventoryEvent {
agent_id: login_response.agent_id,
})
.await
{
Err(CapabilityError {
message: e.to_string(),
})?
}
if let Err(e) = mailbox_addr.send(FetchEnvironmentEvent {}).await {
Err(CapabilityError {
message: e.to_string(),
})?
}
Ok(())
}