Replace callback-based requests/messages with streams

This commit is contained in:
Antonio Scandurra
2021-06-16 14:26:54 +02:00
parent 8b66e0aa7e
commit 8112efd522
8 changed files with 306 additions and 137 deletions
+102 -1
View File
@@ -1,5 +1,8 @@
use crate::rpc_client::{Message, Request, RpcClient};
use postage::prelude::Stream;
use rand::prelude::*;
use std::cmp::Ordering;
use std::{cmp::Ordering, future::Future, sync::Arc};
use zed_rpc::proto;
#[derive(Copy, Clone, Eq, PartialEq, Debug, Hash)]
pub enum Bias {
@@ -53,6 +56,104 @@ where
}
}
pub trait RequestHandler<'a, R: proto::RequestMessage> {
type Output: 'a + Future<Output = anyhow::Result<()>>;
fn handle(
&self,
request: Request<R>,
client: Arc<RpcClient>,
cx: &'a mut gpui::AsyncAppContext,
) -> Self::Output;
}
impl<'a, R, F, Fut> RequestHandler<'a, R> for F
where
R: proto::RequestMessage,
F: Fn(Request<R>, Arc<RpcClient>, &'a mut gpui::AsyncAppContext) -> Fut,
Fut: 'a + Future<Output = anyhow::Result<()>>,
{
type Output = Fut;
fn handle(
&self,
request: Request<R>,
client: Arc<RpcClient>,
cx: &'a mut gpui::AsyncAppContext,
) -> Self::Output {
(self)(request, client, cx)
}
}
pub trait MessageHandler<'a, M: proto::EnvelopedMessage> {
type Output: 'a + Future<Output = anyhow::Result<()>>;
fn handle(
&self,
message: Message<M>,
client: Arc<RpcClient>,
cx: &'a mut gpui::AsyncAppContext,
) -> Self::Output;
}
impl<'a, M, F, Fut> MessageHandler<'a, M> for F
where
M: proto::EnvelopedMessage,
F: Fn(Message<M>, Arc<RpcClient>, &'a mut gpui::AsyncAppContext) -> Fut,
Fut: 'a + Future<Output = anyhow::Result<()>>,
{
type Output = Fut;
fn handle(
&self,
message: Message<M>,
client: Arc<RpcClient>,
cx: &'a mut gpui::AsyncAppContext,
) -> Self::Output {
(self)(message, client, cx)
}
}
pub fn spawn_request_handler<H, R>(
handler: H,
client: &Arc<RpcClient>,
cx: &mut gpui::MutableAppContext,
) where
H: 'static + for<'a> RequestHandler<'a, R>,
R: proto::RequestMessage,
{
let client = client.clone();
let mut requests = smol::block_on(client.add_request_handler::<R>());
cx.spawn(|mut cx| async move {
while let Some(request) = requests.recv().await {
if let Err(err) = handler.handle(request, client.clone(), &mut cx).await {
log::error!("error handling request: {:?}", err);
}
}
})
.detach();
}
pub fn spawn_message_handler<H, M>(
handler: H,
client: &Arc<RpcClient>,
cx: &mut gpui::MutableAppContext,
) where
H: 'static + for<'a> MessageHandler<'a, M>,
M: proto::EnvelopedMessage,
{
let client = client.clone();
let mut messages = smol::block_on(client.add_message_handler::<M>());
cx.spawn(|mut cx| async move {
while let Some(message) = messages.recv().await {
if let Err(err) = handler.handle(message, client.clone(), &mut cx).await {
log::error!("error handling message: {:?}", err);
}
}
})
.detach();
}
pub struct RandomCharIter<T: Rng>(T);
impl<T: Rng> RandomCharIter<T> {