use crate::{constant, resources, routes, squire};
use actix;
use actix_web::{rt, web, Error, HttpRequest, HttpResponse};
use actix_ws::AggregatedMessage;
use fernet::Fernet;
use futures::future;
use futures::stream::StreamExt;
use std::sync::Arc;
use std::time::Duration;
async fn send_system_resources(request: HttpRequest, mut session: actix_ws::Session) {
let host = request.connection_info().host().to_string();
let disk_stats = resources::stream::get_disk_stats();
loop {
let mut system_resources = resources::stream::system_resources();
system_resources.insert("disk_info".to_string(), disk_stats.clone());
let serialized = serde_json::to_string(&system_resources).unwrap();
match session.text(serialized).await {
Ok(_) => (),
Err(err) => {
log::info!("Connection from '{}' has been {}", host, err.to_string().to_lowercase());
break;
}
}
rt::time::sleep(Duration::from_secs(1)).await;
}
}
async fn receive_messages(
mut session: actix_ws::Session,
mut stream: impl futures::Stream<Item=Result<AggregatedMessage, actix_ws::ProtocolError>> + Unpin,
) {
while let Some(msg) = stream.next().await {
match msg {
Ok(AggregatedMessage::Text(text)) => {
session.text(text).await.unwrap();
}
Ok(AggregatedMessage::Binary(bin)) => {
session.binary(bin).await.unwrap();
}
Ok(AggregatedMessage::Ping(msg)) => {
session.pong(&msg).await.unwrap();
}
_ => {}
}
}
}
async fn session_handler(session: actix_ws::Session, duration: i64) {
let session = session.clone();
actix::spawn(async move {
rt::time::sleep(Duration::from_secs(duration as u64)).await;
let _ = session.close(None).await;
});
}
#[route("/ws/system", method = "GET")]
async fn echo(
request: HttpRequest,
fernet: web::Data<Arc<Fernet>>,
session_info: web::Data<Arc<constant::Session>>,
config: web::Data<Arc<squire::settings::Config>>,
stream: web::Payload,
) -> Result<HttpResponse, Error> {
log::info!("Websocket connection initiated");
let auth_response = squire::authenticator::verify_token(&request, &config, &fernet, &session_info);
if !auth_response.ok {
return Ok(routes::auth::failed_auth(auth_response));
}
let (response, session, stream) = match actix_ws::handle(&request, stream) {
Ok(result) => result,
Err(_) => {
return Ok(HttpResponse::ServiceUnavailable().finish());
}
};
let stream = stream
.aggregate_continuations();
rt::spawn(async move {
log::warn!("Connection established");
let send_task = send_system_resources(request.clone(), session.clone());
let receive_task = receive_messages(session.clone(), stream);
let session_task = session_handler(session.clone(), config.session_duration);
future::join3(send_task, receive_task, session_task).await;
});
Ok(response)
}