Working session structure
This commit is contained in:
+24
-19
@@ -1,32 +1,37 @@
|
|||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use session_rs::session::Session;
|
use session_rs::{Method, session::Session, ws::Frame};
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize)]
|
#[derive(Debug, Serialize, Deserialize)]
|
||||||
struct Data {}
|
struct Data {}
|
||||||
|
|
||||||
type Communication = Session<Data, Data, Data, Data, Data>;
|
impl Method for Data {
|
||||||
|
const NAME: &'static str = "data";
|
||||||
|
type Request = ();
|
||||||
|
type Response = ();
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::main(flavor = "current_thread")]
|
#[tokio::main(flavor = "current_thread")]
|
||||||
async fn main() -> session_rs::Result<()> {
|
async fn main() -> session_rs::Result<()> {
|
||||||
let session = Communication::connect("127.0.0.1:8080", "/").await?;
|
let session = Session::connect("127.0.0.1:8080", "/").await?;
|
||||||
|
|
||||||
// Spawn read loop
|
tokio::spawn({
|
||||||
// tokio::spawn({
|
let session = session.clone();
|
||||||
// let session = session.clone();
|
async move {
|
||||||
// async move {
|
loop {
|
||||||
// loop {
|
match session.ws.read().await {
|
||||||
// match session.read().await {
|
Ok(Frame::Text(text)) => {
|
||||||
// Ok(Frame::Text(text)) => {
|
println!("Server says: {}", text);
|
||||||
// println!("Server says: {}", text);
|
}
|
||||||
// }
|
Ok(_) => {}
|
||||||
// Ok(_) => {}
|
Err(_) => break,
|
||||||
// Err(_) => break,
|
}
|
||||||
// }
|
}
|
||||||
// }
|
}
|
||||||
// }
|
});
|
||||||
// });
|
|
||||||
|
|
||||||
session.request(&Data {}).await?;
|
session.request::<Data>(()).await?;
|
||||||
|
|
||||||
|
tokio::time::sleep(tokio::time::Duration::from_millis(1000)).await;
|
||||||
|
|
||||||
session.close().await?;
|
session.close().await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
|
|||||||
+1
-1
@@ -31,7 +31,7 @@ async fn main() -> session_rs::Result<()> {
|
|||||||
println!("Received text: {}", text);
|
println!("Received text: {}", text);
|
||||||
|
|
||||||
// Echo back
|
// Echo back
|
||||||
if let Err(e) = session.send(&serde_json::json!({"echo": text}).to_string()).await {
|
if let Err(e) = session.send(&text).await {
|
||||||
eprintln!("Send error: {:?}", e);
|
eprintln!("Send error: {:?}", e);
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,9 +1,17 @@
|
|||||||
|
use serde::{Deserialize, Serialize};
|
||||||
|
|
||||||
pub mod server;
|
pub mod server;
|
||||||
pub mod session;
|
pub mod session;
|
||||||
pub mod ws;
|
pub mod ws;
|
||||||
|
|
||||||
pub type Result<T> = std::result::Result<T, Error>;
|
pub type Result<T> = std::result::Result<T, Error>;
|
||||||
|
|
||||||
|
pub trait Method {
|
||||||
|
const NAME: &'static str;
|
||||||
|
type Request: Serialize + for<'de> Deserialize<'de>;
|
||||||
|
type Response: Serialize + for<'de> Deserialize<'de>;
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
pub enum Error {
|
pub enum Error {
|
||||||
WebSocket(ws::Error),
|
WebSocket(ws::Error),
|
||||||
|
|||||||
+53
-96
@@ -1,82 +1,45 @@
|
|||||||
use std::{marker::PhantomData, sync::Arc};
|
use std::sync::Arc;
|
||||||
|
|
||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use tokio::sync::Mutex;
|
use tokio::sync::Mutex;
|
||||||
|
|
||||||
use crate::ws::WebSocket;
|
use crate::{Method, ws::WebSocket};
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize)]
|
#[derive(Debug, Serialize, Deserialize)]
|
||||||
pub enum SessionMessageKind {
|
pub enum Message<M: Method> {
|
||||||
Request,
|
Request {
|
||||||
Response,
|
id: u32,
|
||||||
Notification,
|
method: String,
|
||||||
|
data: M::Request,
|
||||||
|
},
|
||||||
|
Response {
|
||||||
|
id: u32,
|
||||||
|
error: bool,
|
||||||
|
result: M::Response,
|
||||||
|
},
|
||||||
|
Notification {
|
||||||
|
method: String,
|
||||||
|
data: M::Request,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Serialize, Deserialize)]
|
pub struct Session {
|
||||||
pub struct SessionMessage<T: Serialize> {
|
|
||||||
id: u32,
|
|
||||||
kind: SessionMessageKind,
|
|
||||||
data: T,
|
|
||||||
}
|
|
||||||
|
|
||||||
pub struct Session<
|
|
||||||
Req: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Res: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerReq: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerRes: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Notification: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
> {
|
|
||||||
_pd: (
|
|
||||||
PhantomData<Req>,
|
|
||||||
PhantomData<Res>,
|
|
||||||
PhantomData<PeerReq>,
|
|
||||||
PhantomData<PeerRes>,
|
|
||||||
PhantomData<Notification>,
|
|
||||||
),
|
|
||||||
pub ws: WebSocket,
|
pub ws: WebSocket,
|
||||||
id: Arc<Mutex<u32>>,
|
id: Arc<Mutex<u32>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<
|
impl Session {
|
||||||
Req: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Res: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerReq: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerRes: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Notification: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
> Session<Req, Res, PeerReq, PeerRes, Notification>
|
|
||||||
{
|
|
||||||
pub fn clone(&self) -> Self {
|
pub fn clone(&self) -> Self {
|
||||||
Self {
|
Self {
|
||||||
_pd: (
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
),
|
|
||||||
ws: self.ws.clone(),
|
ws: self.ws.clone(),
|
||||||
id: self.id.clone(),
|
id: self.id.clone(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<
|
impl Session {
|
||||||
Req: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Res: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerReq: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerRes: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Notification: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
> Session<Req, Res, PeerReq, PeerRes, Notification>
|
|
||||||
{
|
|
||||||
pub fn from_ws(ws: WebSocket) -> Self {
|
pub fn from_ws(ws: WebSocket) -> Self {
|
||||||
Self {
|
Self {
|
||||||
_pd: (
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
PhantomData,
|
|
||||||
),
|
|
||||||
ws,
|
ws,
|
||||||
id: Arc::new(Mutex::new(0)),
|
id: Arc::new(Mutex::new(0)),
|
||||||
}
|
}
|
||||||
@@ -87,55 +50,49 @@ impl<
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<
|
impl Session {
|
||||||
Req: Serialize + for<'a> Deserialize<'a>,
|
pub async fn send<M: Method>(&self, data: &Message<M>) -> crate::Result<()> {
|
||||||
Res: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerReq: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
PeerRes: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
Notification: Serialize + for<'a> Deserialize<'a>,
|
|
||||||
> Session<Req, Res, PeerReq, PeerRes, Notification>
|
|
||||||
{
|
|
||||||
pub async fn send_id<T: Serialize>(
|
|
||||||
&self,
|
|
||||||
id: u32,
|
|
||||||
kind: SessionMessageKind,
|
|
||||||
data: &T,
|
|
||||||
) -> crate::Result<()> {
|
|
||||||
self.ws
|
self.ws
|
||||||
.send_text_payload(&serde_json::to_vec(&SessionMessage { id, kind, data })?)
|
.send_text_payload(&serde_json::to_vec(&data)?)
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn send<T: Serialize>(
|
pub async fn use_id(&self) -> u32 {
|
||||||
&self,
|
let mut id = self.id.lock().await;
|
||||||
kind: SessionMessageKind,
|
*id += 1;
|
||||||
data: &T,
|
*id
|
||||||
) -> crate::Result<()> {
|
}
|
||||||
self.send_id(
|
|
||||||
{
|
pub async fn request<M: Method>(&self, req: M::Request) -> crate::Result<()> {
|
||||||
let mut i = self.id.lock().await;
|
self.send::<M>(&Message::Request {
|
||||||
*i += 1;
|
id: self.use_id().await,
|
||||||
*i
|
method: M::NAME.to_string(),
|
||||||
},
|
data: req,
|
||||||
kind,
|
})
|
||||||
data,
|
|
||||||
)
|
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn request(&self, data: &Req) -> crate::Result<()> {
|
pub async fn respond<M: Method>(
|
||||||
self.send(SessionMessageKind::Request, data).await
|
&self,
|
||||||
|
to: u32,
|
||||||
|
error: bool,
|
||||||
|
res: M::Response,
|
||||||
|
) -> crate::Result<()> {
|
||||||
|
self.send::<M>(&Message::Response {
|
||||||
|
id: to,
|
||||||
|
error,
|
||||||
|
result: res,
|
||||||
|
})
|
||||||
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn respond(&self, to_req: &SessionMessage<PeerReq>, data: &Res) -> crate::Result<()> {
|
pub async fn notify<M: Method>(&self, data: M::Request) -> crate::Result<()> {
|
||||||
self.send_id(to_req.id, SessionMessageKind::Response, data)
|
self.send::<M>(&Message::Notification {
|
||||||
.await
|
method: M::NAME.to_string(),
|
||||||
}
|
data,
|
||||||
|
})
|
||||||
pub async fn notify(&self, data: &Res) -> crate::Result<()> {
|
.await
|
||||||
self.send(SessionMessageKind::Notification, data).await
|
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn close(&self) -> crate::Result<()> {
|
pub async fn close(&self) -> crate::Result<()> {
|
||||||
|
|||||||
Reference in New Issue
Block a user