use actix_web::{rt, web::{Payload, Data, Path, Json}, Error, HttpRequest, HttpResponse, get, post}; //use actix_ws::AggregatedMessage; //use futures::StreamExt as _; #[get("/intercom/.oj_services/{name}")] pub async fn services_ws(req: HttpRequest, stream: Payload, auth: Data, reg: Data, name: Path) -> Result { auth.validate(&req, &format!(".oj_services/{}", urlencoding::encode(&name)))?; let (res, mut session, stream) = actix_ws::handle(&req, stream)?; let mut stream = stream .aggregate_continuations() .max_continuation_size(2_usize.pow(20)); // aggregate continuation frames up to 1MiB let (tx, mut rx) = tokio::sync::mpsc::channel(16); reg.register_web_service(name.clone(), tx.clone()).await; log::debug!("Registered web services intercom websocket for user {}", name); let name_rx: String = name.clone(); let mut session_rx = session.clone(); // start receiver task but don't wait for it rt::spawn(async move { while let Some(msg_result) = stream.recv().await { match msg_result { Ok(msg) => { match msg { actix_ws::AggregatedMessage::Text(s) => { // TODO allow responses log::debug!("Received `{}` ({}) from services websocket receiver for user {}", s, s.len(), name_rx); }, actix_ws::AggregatedMessage::Binary(_) => {}, actix_ws::AggregatedMessage::Ping(p) => { log::debug!("Received {}B ping from services websocket for user {}: {}", p.len(), name_rx, String::from_utf8_lossy(&p)); session_rx.pong(&p).await.unwrap_or_default(); }, actix_ws::AggregatedMessage::Pong(p) => { log::debug!("Received {}B pong from services websocket for user {}: {}", p.len(), name_rx, String::from_utf8_lossy(&p)); }, actix_ws::AggregatedMessage::Close(reason) => { if let Some(reason) = reason { log::info!("Closed services websocket receiver for user {} with reason {:?}", name_rx, reason); } else { log::warn!("Closed services websocket receiver for user {} without reason", name_rx); } break; }, } }, Err(e) => { log::warn!("Received services websocket error for user {}: {}", name_rx, e); }, } } tx.send(crate::robocraft::intercom::IntercomOp::Info(crate::robocraft::intercom::IntercomInfo::Close)).await.unwrap_or_default(); }); // start sender task but don't wait for it rt::spawn(async move { let mut is_ok = false; session.ping(b"openjam").await.unwrap_or_default(); while let Some(op) = rx.recv().await { match op { super::IntercomOp::Message(msg) => { if let Err(e) = session.text(serde_json::to_string(&msg).unwrap()).await { log::warn!("Failed to send services intercom to user {}: {}", name, e); break; } }, super::IntercomOp::Info(info) => { match info { super::IntercomInfo::Close => { is_ok = true; break; }, } } } } if !is_ok { reg.remove_web_service(name.clone()).await; } rx.close(); session.close(Some(actix_ws::CloseReason { code: actix_ws::CloseCode::Normal, description: Some("End of channel".to_owned()), })).await.expect("Failed to close a services intercom websocket session"); log::debug!("Web services intercom websocket closed for {}", name); }); // respond immediately with response connected to WS session Ok(res) } #[post("/intercom/.oj_services/{name}/messages")] pub async fn service_msg(req: HttpRequest, body: Json, auth: Data, reg: Data, name: Path) -> Result { log::debug!("Got intercom message from {} to {:?}", name, body.public_ids.as_slice()); auth.validate(&req, &format!(".oj_services/{}/messages", name))?; log::debug!("Authenticated intercom message from {} to {:?}", name, body.public_ids.as_slice()); reg.broadcast_web_service_message(body.0).await; Ok(HttpResponse::NoContent().finish()) }