use std::convert::TryFrom;
use std::fmt::Debug;
use futures::{Stream, StreamExt};
use tonic::transport::Channel;
use tracing::{instrument, trace};
use crate::data::completion::DamlCompletionResponse;
use crate::data::offset::DamlLedgerOffset;
use crate::data::{DamlError, DamlResult};
use crate::grpc_protobuf::com::daml::ledger::api::v2::CompletionStreamRequest;
use crate::grpc_protobuf::com::daml::ledger::api::v2::command_completion_service_client::CommandCompletionServiceClient;
use crate::service::common::make_request;
#[derive(Debug)]
pub struct DamlCommandCompletionService<'a> {
channel: Channel,
auth_token: Option<&'a str>,
}
impl<'a> DamlCommandCompletionService<'a> {
pub fn new(channel: Channel, auth_token: Option<&'a str>) -> Self {
Self {
channel,
auth_token,
}
}
pub fn with_token(self, auth_token: &'a str) -> Self {
Self {
auth_token: Some(auth_token),
..self
}
}
#[instrument(skip(self))]
pub async fn get_completion_stream(
&self,
user_id: impl Into<String> + Debug,
parties: impl Into<Vec<String>> + Debug,
begin_exclusive: impl Into<DamlLedgerOffset> + Debug,
) -> DamlResult<impl Stream<Item = DamlResult<DamlCompletionResponse>>> {
let payload = CompletionStreamRequest {
user_id: user_id.into(),
parties: parties.into(),
begin_exclusive: begin_exclusive.into().value(),
};
trace!(payload = ?payload, token = ?self.auth_token);
let completion_stream =
self.client().completion_stream(make_request(payload, self.auth_token)?).await?.into_inner();
Ok(completion_stream.inspect(|response| trace!(?response)).map(|item| match item {
Ok(completion) => DamlCompletionResponse::try_from(completion),
Err(e) => Err(DamlError::from(e)),
}))
}
fn client(&self) -> CommandCompletionServiceClient<Channel> {
CommandCompletionServiceClient::new(self.channel.clone())
}
}