use std::time::Duration; use crate::relay_loop::run_relay_loop; use crate::settings::Settings; use emgauwa_lib::constants::WEBSOCKET_RETRY_TIMEOUT; use emgauwa_lib::db::{DbController, DbRelay}; use emgauwa_lib::handlers::v1::ws::controllers::ControllerWsAction; use emgauwa_lib::models::{Controller, FromDbModel}; use emgauwa_lib::types::ControllerUid; use emgauwa_lib::{db, utils}; use futures::{SinkExt, StreamExt}; use sqlx::pool::PoolConnection; use sqlx::Sqlite; use tokio::time; use tokio_tungstenite::tungstenite::Error; use tokio_tungstenite::{connect_async, tungstenite::protocol::Message}; use utils::init_logging; mod driver; mod relay_loop; mod settings; async fn create_this_controller( conn: &mut PoolConnection, settings: &Settings, ) -> DbController { DbController::create( conn, &ControllerUid::default(), &settings.name, i64::try_from(settings.relays.len()).expect("Too many relays"), ) .await .expect("Failed to create controller") } async fn create_this_relay( conn: &mut PoolConnection, this_controller: &DbController, settings_relay: &settings::Relay, ) -> DbRelay { DbRelay::create( conn, &settings_relay.name, settings_relay.number.unwrap(), this_controller, ) .await .expect("Failed to create relay") } #[tokio::main] async fn main() { let settings = settings::init(); init_logging(&settings.logging.level); let pool = db::init(&settings.database).await; let mut conn = pool.acquire().await.unwrap(); let db_controller = DbController::get_all(&mut conn) .await .expect("Failed to get controller from database") .pop() .unwrap_or_else(|| { futures::executor::block_on(create_this_controller(&mut conn, &settings)) }); for relay in &settings.relays { if DbRelay::get_by_controller_and_num(&mut conn, &db_controller, relay.number.unwrap()) .await .expect("Failed to get relay from database") .is_none() { create_this_relay(&mut conn, &db_controller, relay).await; } } let db_controller = db_controller .update(&mut conn, &db_controller.name, settings.relays.len() as i64) .await .unwrap(); let this = Controller::from_db_model(&mut conn, db_controller) .expect("Failed to convert database models"); let url = format!( "ws://{}:{}/api/v1/ws/controllers", settings.core.host, settings.core.port ); tokio::spawn(run_relay_loop(settings)); loop { time::sleep(WEBSOCKET_RETRY_TIMEOUT).await; let connect_result = connect_async(&url).await; if let Err(err) = connect_result { log::warn!( "Failed to connect to websocket: {}. Retrying in {} seconds...", err, WEBSOCKET_RETRY_TIMEOUT.as_secs() ); continue; } let (ws_stream, _) = connect_result.unwrap(); let (mut write, read) = ws_stream.split(); let ws_action = ControllerWsAction::Register(this.clone()); let ws_action_json = serde_json::to_string(&ws_action).unwrap(); write.send(Message::text(ws_action_json)).await.unwrap(); let read_handler = read.for_each(handle_message); read_handler.await; log::warn!( "Lost connection to websocket. Retrying in {} seconds...", WEBSOCKET_RETRY_TIMEOUT.as_secs() ); } } pub async fn handle_message(message_result: Result) { match message_result { Ok(message) => { if let Message::Text(msg_text) = message { log::debug!("{}", msg_text) } } Err(err) => log::debug!("Error: {}", err), } }