use crate::animation::{AnimationPath, AnimationQueue, scene_instance_ready, update_animations};
use crate::chat::ChatPlugin;
use crate::environment::{LandPlugin, LandUpdateEvent};
use crate::errors::{NotLoggedIn, PacketSendError, PortError, ShareDirError};
use crate::login;
use crate::render::{
AgentIDMap, MeshQueue, MeshUpdateEvent, SceneIDMap, follow_gltf_with_offset,
handle_camera_update, handle_mesh_update, render_land, render_meshes,
};
use crate::sky::SkyboxUpdateEvent;
use crate::subscriber::listen_for_core_events;
use crate::water::{WaterPlugin, WaterUpdateEvent};
use actix_rt::System;
use benthic_protocol::messages::ui::agent_update::AgentUpdate;
use benthic_protocol::messages::ui::camera_position::CameraPosition;
use benthic_protocol::messages::ui::coarse_location_update::CoarseLocationUpdate;
use benthic_protocol::messages::ui::errors::SessionError;
use benthic_protocol::messages::ui::login_error::LoginError;
use benthic_protocol::messages::ui::login_response::LoginResponse;
use benthic_protocol::messages::ui::play_animation::PlayAnimation;
use benthic_protocol::messages::ui::ui_messages::{UIMessage, UIResponse};
use benthic_protocol::session::initialize_share_dir;
use bevy::animation::AnimatedBy;
use bevy::app::App;
use bevy::camera::visibility::DynamicSkinnedMeshBounds;
use bevy::mesh::skinning::SkinnedMesh;
use bevy::platform::collections::HashMap;
use bevy::prelude::*;
use bevy::tasks::AsyncComputeTaskPool;
use bevy::window::WindowCloseRequested;
use bevy_gltf::{Gltf, GltfMaterialName, GltfMeshName, GltfSceneName};
use crossbeam_channel::{Receiver, Sender, unbounded};
use metaverse_core::initialize::initialize;
use portpicker::pick_unused_port;
use std::net::UdpSocket;
use std::path::PathBuf;
pub const VIEWER_NAME: &str = "benthic";
#[derive(Resource)]
pub struct SessionData {
pub login_response: Option<LoginResponse>,
pub avatar_location: Vec3,
}
#[derive(Clone, Copy, Default, Eq, PartialEq, Debug, Hash, States)]
pub enum ViewerState {
#[default]
Login,
Loading,
Chat,
}
#[derive(Resource)]
pub struct Sockets {
pub ui_to_core_socket: u16,
pub core_to_ui_socket: u16,
}
#[derive(Resource)]
struct EventChannel {
pub sender: Sender<UIMessage>,
pub receiver: Receiver<UIMessage>,
}
#[derive(Resource)]
pub struct ChatMessages {
pub messages: Vec<ChatFromClientMessage>,
}
#[derive(Resource)]
pub struct ShareDir {
pub _path: PathBuf,
pub login_cred_path: PathBuf,
}
pub struct ChatFromClientMessage {
pub user: String,
pub message: String,
}
#[derive(Message)]
pub struct LoginResponseEvent {
pub value: Result<LoginResponse, LoginError>,
}
#[derive(Message)]
pub struct CameraUpdateEvent {
pub value: CameraPosition,
}
#[derive(Message, Clone)]
pub struct PlayAnimationEvent {
pub value: PlayAnimation,
}
#[derive(Message)]
pub struct CoarseLocationUpdateEvent {
pub _value: CoarseLocationUpdate,
}
#[derive(Resource)]
pub struct AgentUpdateTimer(Timer);
#[derive(Message)]
struct DisableSimulatorEvent;
#[derive(Message)]
struct LogoutRequestEvent;
fn setup_share_dir() -> Result<PathBuf, ShareDirError> {
Ok(initialize_share_dir()?)
}
pub struct MetaversePlugin;
impl Plugin for MetaversePlugin {
fn build(&self, app: &mut App) {
let (s1, r1) = unbounded();
let ui_to_core_socket =
pick_unused_port().unwrap_or_else(|| panic!("{:?}", PortError::PortPickerError()));
let core_to_ui_socket =
pick_unused_port().unwrap_or_else(|| panic!("{:?}", PortError::PortPickerError()));
let local_share_dir = match setup_share_dir() {
Err(e) => {
panic!("{:?}", e)
}
Ok(path) => {
info!("Created directory: {:?}", path);
path
}
};
let login_cred_path = local_share_dir.join("login_conf.json");
let login_data = match login::load_login_data(&login_cred_path) {
Ok(login_data) => login_data,
Err(e) => {
error!("{:?}", e.to_string());
login::LoginData::default()
}
};
app.init_state::<ViewerState>()
.add_plugins(WaterPlugin)
.add_plugins(LandPlugin)
.add_plugins(ChatPlugin)
.insert_resource(ClearColor(Color::BLACK))
.insert_resource(SessionData {
login_response: None,
avatar_location: Vec3::ZERO,
})
.insert_resource(Sockets {
ui_to_core_socket,
core_to_ui_socket,
})
.insert_resource(login_data)
.insert_resource(EventChannel {
sender: s1,
receiver: r1,
})
.insert_resource(ShareDir {
_path: local_share_dir,
login_cred_path,
})
.insert_resource(AgentIDMap {
entities: HashMap::new(),
})
.insert_resource(SceneIDMap {
entities: HashMap::new(),
})
.insert_resource(AnimationQueue {
pending: HashMap::new(),
})
.insert_resource(MeshQueue { pending: vec![] })
.add_message::<LoginResponseEvent>()
.add_message::<CameraUpdateEvent>()
.add_message::<CoarseLocationUpdateEvent>()
.add_message::<MeshUpdateEvent>()
.add_message::<LandUpdateEvent>()
.add_message::<DisableSimulatorEvent>()
.add_message::<LogoutRequestEvent>()
.add_message::<SkyboxUpdateEvent>()
.register_type::<GltfSceneName>()
.register_type::<Transform>()
.register_type::<GlobalTransform>()
.register_type::<TransformTreeChanged>()
.register_type::<Children>()
.register_type::<Visibility>()
.register_type::<ChildOf>()
.register_type::<InheritedVisibility>()
.register_type::<ViewVisibility>()
.register_type::<AnimationPlayer>()
.register_type::<Name>()
.register_type::<Mesh3d>()
.register_type::<bevy::camera::primitives::Aabb>()
.register_type::<SkinnedMesh>()
.register_type::<GltfMeshName>()
.register_type::<GltfMaterialName>()
.register_type::<AnimatedBy>()
.register_type::<DynamicSkinnedMeshBounds>()
.add_systems(Startup, start_listener)
.add_systems(Startup, setup_timers)
.add_systems(Startup, start_core)
.add_systems(Update, handle_window_close)
.add_systems(Update, handle_logout)
.add_systems(Update, handle_queue)
.add_systems(Update, handle_login_response)
.add_systems(Update, handle_disconnect)
.add_systems(Update, handle_mesh_update)
.add_systems(Update, update_animations)
.add_systems(Update, handle_camera_update)
.add_systems(Update, follow_gltf_with_offset)
.add_systems(Update, render_land)
.add_systems(Update, render_meshes)
.add_systems(
Update,
send_agent_update.run_if(in_state(ViewerState::Chat)),
)
.add_observer(scene_instance_ready);
}
}
pub fn handle_login_response(
mut ev_loginresponse: MessageReader<LoginResponseEvent>,
mut viewer_state: ResMut<NextState<ViewerState>>,
mut session_data: ResMut<SessionData>,
) {
for response in ev_loginresponse.read() {
match response.value.as_ref() {
Ok(login_response) => {
viewer_state.set(ViewerState::Chat);
session_data.login_response = Some(login_response.clone());
}
Err(_) => viewer_state.set(ViewerState::Login),
}
}
}
pub fn send_packet_to_core(packet: &[u8], sockets: &Res<Sockets>) -> Result<(), PacketSendError> {
let client_socket = UdpSocket::bind("0.0.0.0:0")?;
client_socket.send_to(packet, format!("127.0.0.1:{}", sockets.ui_to_core_socket))?;
Ok(())
}
fn handle_window_close(
mut events: MessageReader<WindowCloseRequested>,
mut exit: MessageWriter<AppExit>,
mut logout: MessageWriter<LogoutRequestEvent>,
viewer_state: Res<State<ViewerState>>,
) {
for _ in events.read() {
info!("Window close requested. Exiting...");
if *viewer_state == ViewerState::Chat {
info!("Sending Logout");
logout.write(LogoutRequestEvent);
}
exit.write(AppExit::Success);
}
}
fn setup_timers(mut commands: Commands) {
commands.insert_resource(AgentUpdateTimer(Timer::from_seconds(
0.1,
TimerMode::Repeating,
)));
}
fn handle_disconnect(
mut ev_disable_simulator: MessageReader<DisableSimulatorEvent>,
mut viewer_state: ResMut<NextState<ViewerState>>,
) {
for _ in ev_disable_simulator.read() {
viewer_state.set(ViewerState::Login);
}
}
#[allow(clippy::all)]
fn handle_queue(
event_channel: Res<EventChannel>,
mut ev_loginresponse: MessageWriter<LoginResponseEvent>,
mut ev_coarselocationupdate: MessageWriter<CoarseLocationUpdateEvent>,
mut ev_disable_simulator: MessageWriter<DisableSimulatorEvent>,
mut ev_mesh_update: MessageWriter<MeshUpdateEvent>,
mut ev_land_update: MessageWriter<LandUpdateEvent>,
mut ev_camera_update: MessageWriter<CameraUpdateEvent>,
mut ev_water_update: MessageWriter<WaterUpdateEvent>,
mut ev_skybox_update: MessageWriter<SkyboxUpdateEvent>,
mut chat_messages: ResMut<ChatMessages>,
mut animation_queue: ResMut<AnimationQueue>,
asset_server: Res<AssetServer>,
) {
let receiver = event_channel.receiver.clone();
while let Ok(event) = receiver.try_recv() {
match event {
UIMessage::LandUpdate(land_update) => {
ev_land_update.write(LandUpdateEvent { value: land_update });
}
UIMessage::LoginResponse(login_response) => {
ev_loginresponse.write(LoginResponseEvent {
value: Ok(login_response),
});
}
UIMessage::MeshUpdate(mesh_update) => {
ev_mesh_update.write(MeshUpdateEvent { value: mesh_update });
}
UIMessage::PlayAnimation(play_animation) => {
let gltf_handle: Handle<Gltf> =
asset_server.load(play_animation.animation_path.clone());
animation_queue.pending.insert(
play_animation.player_id,
AnimationPath {
path_on_disk: play_animation.animation_path.clone(),
gltf_handle,
},
);
}
UIMessage::CoarseLocationUpdate(coarse_location_update) => {
ev_coarselocationupdate.write(CoarseLocationUpdateEvent {
_value: coarse_location_update,
});
}
UIMessage::ChatFromSimulator(chat_from_simulator) => {
chat_messages.messages.push(ChatFromClientMessage {
user: chat_from_simulator.from_name,
message: chat_from_simulator.message,
});
}
UIMessage::DisableSimulator(_) => {
ev_disable_simulator.write(DisableSimulatorEvent {});
}
UIMessage::CameraPosition(data) => {
ev_camera_update.write(CameraUpdateEvent { value: data });
}
UIMessage::WaterUpdate(data) => {
ev_water_update.write(WaterUpdateEvent { value: data });
}
UIMessage::SkyboxUpdate(data) => {
ev_skybox_update.write(SkyboxUpdateEvent { value: data });
}
UIMessage::Error(error) => match error {
SessionError::Login(e) => {
ev_loginresponse.write(LoginResponseEvent {
value: Err(e.clone()),
});
error!("{:?}", e)
}
SessionError::MailboxSession(e) => {
info!("MailboxError {:?}", e)
}
SessionError::AckError(e) => {
info!("AckError {:?}", e)
}
SessionError::CircuitCode(e) => {
info!("CircuitcodeError {:?}", e)
}
SessionError::CompleteAgentMovement(e) => {
info!("CompleteAgentMovmentError {:?}", e)
}
SessionError::Capability(e) => {
info!("CapabilityError {:?}", e)
}
SessionError::IOError(e) => {
info!("IOError {:?}", e)
}
SessionError::FeatureError(e) => {
info!("FeatureError {:?}", e)
}
},
};
}
}
fn start_listener(sockets: Res<Sockets>, event_queue: Res<EventChannel>) {
let outgoing_socket = sockets.core_to_ui_socket;
let thread_pool = AsyncComputeTaskPool::get();
let sender = event_queue.sender.clone();
thread_pool
.spawn(async move {
listen_for_core_events(format!("127.0.0.1:{}", outgoing_socket), sender).await
})
.detach();
}
fn start_core(sockets: Res<Sockets>) {
let server_to_ui_socket = sockets.core_to_ui_socket;
let ui_to_server_socket = sockets.ui_to_core_socket;
std::thread::spawn(move || {
System::new().block_on(async {
match initialize(ui_to_server_socket, server_to_ui_socket).await {
Ok(handle) => {
match handle.await {
Ok(()) => info!("Listener exited successfully!"),
Err(e) => error!("Listener exited with error {:?}", e),
};
}
Err(err) => {
error!("Failed to start client: {:?}", err);
}
}
});
});
}
pub fn retrieve_login_response<'a>(
session_data: &'a Res<SessionData>,
viewer_state: &mut ResMut<NextState<ViewerState>>,
) -> Result<&'a LoginResponse, NotLoggedIn> {
match &session_data.login_response {
Some(session) => Ok(session), None => {
viewer_state.set(ViewerState::Login);
error!("{:?}", NotLoggedIn::NotLoggedInError());
Err(NotLoggedIn::NotLoggedInError())
}
}
}
fn send_agent_update(
sockets: Res<Sockets>,
time: Res<Time>,
session: Res<SessionData>,
mut timer: ResMut<AgentUpdateTimer>,
) {
if !timer.0.tick(time.delta()).just_finished()
&& let Err(e) = send_packet_to_core(
&UIResponse::new_agent_update(AgentUpdate {
camera_center: session.avatar_location,
..Default::default()
})
.to_bytes(),
&sockets,
)
{
error!("{:?}", e)
};
}
fn handle_logout(mut events: MessageReader<LogoutRequestEvent>, sockets: Res<Sockets>) {
for _ in events.read() {
if let Err(e) = send_packet_to_core(&UIResponse::new_logout().to_bytes(), &sockets) {
error!("{:?}", e)
};
}
}