Change the context API to make Context-Ownership more flexible
This commit is contained in:
@@ -6,7 +6,7 @@ edition = "2021"
|
||||
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
|
||||
|
||||
[dependencies]
|
||||
int-enum = { version = "0.5", features = ["serde", "convert"] }
|
||||
int-enum = { version = "1.2.0" }
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
serde_json = "1"
|
||||
nanoid = "0.4"
|
||||
|
||||
@@ -112,19 +112,27 @@ pub trait JRPCServerService: Send + Sync {
|
||||
&self,
|
||||
request: &JRPCRequest,
|
||||
function: &str,
|
||||
ctx: Self::Context,
|
||||
ctx: &Self::Context,
|
||||
) -> Result<(bool, Value)>;
|
||||
}
|
||||
|
||||
pub type JRPCServiceHandle<Context> = Arc<dyn JRPCServerService<Context = Context>>;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct JRPCSession<Context> {
|
||||
server: JRPCServer<Context>,
|
||||
message_sender: Sender<JRPCResult>,
|
||||
}
|
||||
|
||||
impl<Context: Clone + Send + Sync + 'static> JRPCSession<Context> {
|
||||
impl<Context> Clone for JRPCSession<Context> {
|
||||
fn clone(&self) -> Self {
|
||||
JRPCSession {
|
||||
server: self.server.clone(),
|
||||
message_sender: self.message_sender.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<Context: Send + Sync + 'static> JRPCSession<Context> {
|
||||
pub fn new(server: JRPCServer<Context>, sender: Sender<JRPCResult>) -> Self {
|
||||
JRPCSession {
|
||||
server,
|
||||
@@ -157,64 +165,69 @@ impl<Context: Clone + Send + Sync + 'static> JRPCSession<Context> {
|
||||
pub fn handle_request(&self, request: JRPCRequest, ctx: Context) -> () {
|
||||
let session = self.clone();
|
||||
tokio::task::spawn(async move {
|
||||
info!("Received request: {}", request.method);
|
||||
trace!("Request data: {:?}", request);
|
||||
let method: Vec<&str> = request.method.split('.').collect();
|
||||
if method.len() != 2 {
|
||||
warn!("Invalid method received: {}", request.method);
|
||||
return;
|
||||
}
|
||||
let service = method[0];
|
||||
let function = method[1];
|
||||
let context = ctx;
|
||||
session.handle_request_awaiting(request, &context).await;
|
||||
});
|
||||
}
|
||||
|
||||
let service = session.server.services.get(service);
|
||||
if let Some(service) = service {
|
||||
let result = service.handle(&request, function, ctx).await;
|
||||
match result {
|
||||
Ok((is_send, result)) => {
|
||||
if is_send && request.id.is_some() {
|
||||
let result = session
|
||||
.message_sender
|
||||
.send(JRPCResult {
|
||||
jsonrpc: "2.0".to_string(),
|
||||
id: request.id.unwrap(),
|
||||
result: Some(result),
|
||||
error: None,
|
||||
})
|
||||
.await;
|
||||
if let Err(err) = result {
|
||||
warn!("Error while sending result: {}", err);
|
||||
}
|
||||
pub async fn handle_request_awaiting(&self, request: JRPCRequest, ctx: &Context) -> () {
|
||||
info!("Received request: {}", request.method);
|
||||
trace!("Request data: {:?}", request);
|
||||
let method: Vec<&str> = request.method.split('.').collect();
|
||||
if method.len() != 2 {
|
||||
warn!("Invalid method received: {}", request.method);
|
||||
return;
|
||||
}
|
||||
let service = method[0];
|
||||
let function = method[1];
|
||||
|
||||
let service = self.server.services.get(service);
|
||||
if let Some(service) = service {
|
||||
let result = service.handle(&request, function, ctx).await;
|
||||
match result {
|
||||
Ok((is_send, result)) => {
|
||||
if is_send && request.id.is_some() {
|
||||
let result = self
|
||||
.message_sender
|
||||
.send(JRPCResult {
|
||||
jsonrpc: "2.0".to_string(),
|
||||
id: request.id.unwrap(),
|
||||
result: Some(result),
|
||||
error: None,
|
||||
})
|
||||
.await;
|
||||
if let Err(err) = result {
|
||||
warn!("Error while sending result: {}", err);
|
||||
}
|
||||
}
|
||||
Err(err) => {
|
||||
warn!("Error while handling request: {}", err);
|
||||
session
|
||||
.send_error(
|
||||
request,
|
||||
format!("Error while handling request: {}", err),
|
||||
1,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
warn!("Service not found: {}", method[0]);
|
||||
session
|
||||
.send_error(request, "Service not found".to_string(), 1)
|
||||
.await;
|
||||
return;
|
||||
Err(err) => {
|
||||
warn!("Error while handling request: {}", err);
|
||||
self.send_error(request, format!("Error while handling request: {}", err), 1)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
});
|
||||
} else {
|
||||
warn!("Service not found: {}", method[0]);
|
||||
self.send_error(request, "Service not found".to_string(), 1)
|
||||
.await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct JRPCServer<Context> {
|
||||
services: HashMap<String, JRPCServiceHandle<Context>>,
|
||||
}
|
||||
|
||||
impl<Context: Clone + Send + Sync + 'static> JRPCServer<Context> {
|
||||
impl<Context> Clone for JRPCServer<Context> {
|
||||
fn clone(&self) -> Self {
|
||||
JRPCServer {
|
||||
services: self.services.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<Context: Send + Sync + 'static> JRPCServer<Context> {
|
||||
pub fn new() -> Self {
|
||||
JRPCServer {
|
||||
services: HashMap::new(),
|
||||
|
||||
Reference in New Issue
Block a user