use anyhow::Result;
use tokio::sync::mpsc::Sender;
use tokio_stream::wrappers::ReceiverStream;
use tokio_util::sync::CancellationToken;
use tonic::Streaming;
use zelos_proto::trace::{
subscribe_request, trace_subscribe_client, SubscribeCommand, SubscribeRequest,
SubscribeResponse, UnsubscribeCommand,
};
pub struct TraceSubscribeClient {
req_sender: Sender<SubscribeRequest>,
}
impl TraceSubscribeClient {
pub async fn new(
sender: zelos_trace_types::ipc::Sender,
cancellation_token: CancellationToken,
address: String,
) -> Result<(Self, impl Future<Output = Result<()>>)> {
let channel = zelos_proto::channel::create_channel(address)?;
let mut client = trace_subscribe_client::TraceSubscribeClient::new(channel)
.max_decoding_message_size(zelos_proto::MAX_GRPC_MESSAGE_SIZE)
.max_encoding_message_size(zelos_proto::MAX_GRPC_MESSAGE_SIZE);
let (req_sender, req_receiver) = tokio::sync::mpsc::channel(1);
let request_stream = tonic::Request::new(ReceiverStream::new(req_receiver));
let resp = client.subscribe(request_stream).await?;
let future = Self::run(resp.into_inner(), sender.clone(), cancellation_token);
Ok((Self { req_sender }, future))
}
async fn run(
mut stream: Streaming<SubscribeResponse>,
sender: zelos_trace_types::ipc::Sender,
cancellation_token: CancellationToken,
) -> Result<()> {
loop {
tokio::select! {
msg = stream.message() => {
match msg {
Ok(Some(response)) => {
let ipc = response.as_ipc()?;
for m in ipc {
sender.send_async(m).await?;
}
}
Ok(None) => {
return Ok(());
}
Err(e) => {
return Err(e.into());
}
}
}
_ = cancellation_token.cancelled() => return Ok(())
}
}
}
pub async fn subscribe(&self, filter: Option<String>, start_time: Option<i64>) -> Result<()> {
self.req_sender
.send(SubscribeRequest {
cmd: Some(subscribe_request::Cmd::Subscribe(SubscribeCommand {
filter,
start_time,
})),
})
.await?;
Ok(())
}
pub async fn unsubscribe(&self, filter: Option<String>) -> Result<()> {
self.req_sender
.send(SubscribeRequest {
cmd: Some(subscribe_request::Cmd::Unsubscribe(UnsubscribeCommand {
filter,
})),
})
.await?;
Ok(())
}
pub async fn subscribe_all(&self) -> Result<()> {
self.subscribe(None, None).await
}
pub async fn unsubscribe_all(&self) -> Result<()> {
self.unsubscribe(None).await
}
}