use crate::dapr::proto::common::v1::JobFailurePolicyConstant;
use crate::dapr::proto::common::v1::job_failure_policy::Policy;
use crate::dapr::proto::{common::v1 as common_v1, runtime::v1 as dapr_v1};
use crate::error::Error;
#[cfg(feature = "workflow")]
use crate::workflow;
use async_trait::async_trait;
use futures::StreamExt;
use prost_types::Any;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
use std::time::Duration;
use tokio::io::AsyncRead;
use tonic::codegen::tokio_stream;
use tonic::service::interceptor::InterceptedService;
use tonic::{Request, transport::Channel as TonicChannel};
use tonic::{Status, Streaming};
pub mod config;
pub mod interceptor;
pub use config::{
API_TOKEN_METADATA_KEY, APP_API_TOKEN_ENV, ClientOptions, DAPR_API_TOKEN_ENV,
DAPR_CLIENT_TIMEOUT_SECONDS_ENV, DAPR_GRPC_ENDPOINT_ENV, DAPR_GRPC_PORT_ENV,
DEFAULT_CLIENT_TIMEOUT_SECONDS, DEFAULT_DAPR_GRPC_PORT, default_sidecar_address,
};
pub use interceptor::{ApiTokenInterceptor, AppApiTokenLayer, AppApiTokenService};
#[derive(Clone)]
pub struct Client<T>(T, String);
impl<T: DaprInterface> Client<T> {
#[deprecated(
since = "0.19.0",
note = "Will be removed in 0.20.0. Use Client::new() or Client::from_options()."
)]
pub async fn connect(addr: String) -> Result<Self, Error> {
let port: u16 = std::env::var("DAPR_GRPC_PORT")?.parse()?;
let address = format!("{addr}:{port}");
Ok(Client(T::connect(address.clone()).await?, address))
}
#[deprecated(
since = "0.19.0",
note = "Will be removed in 0.20.0. Use Client::new() or Client::from_options()."
)]
pub async fn connect_with_port(addr: String, port: String) -> Result<Self, Error> {
let port: u16 = match port.parse::<u16>() {
Ok(p) => p,
Err(_) => {
panic!("Port must be a number between 1 and 65535");
}
};
let address = format!("{addr}:{port}");
Ok(Client(T::connect(address.clone()).await?, address))
}
#[cfg(feature = "workflow")]
pub async fn new_workflow_client(&self) -> workflow::Result<workflow::WorkflowClient> {
workflow::WorkflowClient::new_with_address(self.1.clone()).await
}
pub async fn invoke_service<I, M>(
&mut self,
app_id: I,
method_name: M,
data: Option<Any>,
) -> Result<InvokeServiceResponse, Error>
where
I: Into<String>,
M: Into<String>,
{
self.0
.invoke_service(InvokeServiceRequest {
id: app_id.into(),
message: common_v1::InvokeRequest {
method: method_name.into(),
data,
..Default::default()
}
.into(),
})
.await
}
pub async fn invoke_binding<S>(
&mut self,
name: S,
data: Vec<u8>,
operation: S,
metadata: Option<HashMap<String, String>>,
) -> Result<InvokeBindingResponse, Error>
where
S: Into<String>,
{
self.0
.invoke_binding(InvokeBindingRequest {
name: name.into(),
data,
operation: operation.into(),
metadata: metadata.unwrap_or_default(),
})
.await
}
pub async fn invoke_output_binding<S>(&mut self, name: S, operation: S) -> Result<(), Error>
where
S: Into<String>,
{
self.0
.invoke_binding(InvokeBindingRequest {
name: name.into(),
operation: operation.into(),
..Default::default()
})
.await
.map(|_| ())
}
pub async fn publish_event<S>(
&mut self,
pubsub_name: S,
topic: S,
data_content_type: S,
data: Vec<u8>,
metadata: Option<HashMap<String, String>>,
) -> Result<(), Error>
where
S: Into<String>,
{
let mut mdata = HashMap::<String, String>::new();
if let Some(m) = metadata {
mdata = m;
}
self.0
.publish_event(PublishEventRequest {
pubsub_name: pubsub_name.into(),
topic: topic.into(),
data_content_type: data_content_type.into(),
data,
metadata: mdata,
})
.await
}
pub async fn get_secret<S>(&mut self, store_name: S, key: S) -> Result<GetSecretResponse, Error>
where
S: Into<String>,
{
self.0
.get_secret(GetSecretRequest {
store_name: store_name.into(),
key: key.into(),
..Default::default()
})
.await
}
pub async fn get_bulk_secret<S>(
&mut self,
store_name: S,
metadata: Option<HashMap<String, String>>,
) -> Result<GetBulkSecretResponse, Error>
where
S: Into<String>,
{
self.0
.get_bulk_secret(GetBulkSecretRequest {
store_name: store_name.into(),
metadata: metadata.unwrap_or_default(),
})
.await
}
pub async fn get_state<S>(
&mut self,
store_name: S,
key: S,
metadata: Option<HashMap<String, String>>,
) -> Result<GetStateResponse, Error>
where
S: Into<String>,
{
let mut mdata = HashMap::<String, String>::new();
if let Some(m) = metadata {
mdata = m;
}
self.0
.get_state(GetStateRequest {
store_name: store_name.into(),
key: key.into(),
metadata: mdata,
..Default::default()
})
.await
}
pub async fn save_state<S>(
&mut self,
store_name: S,
key: S,
value: Vec<u8>,
etag: Option<Etag>,
metadata: Option<HashMap<String, String>>,
options: Option<StateOptions>,
) -> Result<(), Error>
where
S: Into<String>,
{
let states = vec![StateItem {
key: key.into(),
value,
etag,
metadata: metadata.unwrap_or_default(),
options,
}];
self.save_bulk_states(store_name, states).await
}
pub async fn save_bulk_states<S, I>(&mut self, store_name: S, items: I) -> Result<(), Error>
where
S: Into<String>,
I: Into<Vec<StateItem>>,
{
self.0
.save_state(SaveStateRequest {
store_name: store_name.into(),
states: items.into(),
})
.await
}
pub async fn query_state_alpha1<S>(
&mut self,
store_name: S,
query: Value,
metadata: Option<HashMap<String, String>>,
) -> Result<QueryStateResponse, Error>
where
S: Into<String>,
{
let mut mdata = HashMap::<String, String>::new();
if let Some(m) = metadata {
mdata = m;
}
self.0
.query_state_alpha1(QueryStateRequest {
store_name: store_name.into(),
query: serde_json::to_string(&query).unwrap(),
metadata: mdata,
})
.await
}
pub async fn delete_bulk_state<I, K>(&mut self, store_name: K, states: I) -> Result<(), Error>
where
I: IntoIterator<Item = (K, Vec<u8>)>,
K: Into<String>,
{
self.0
.delete_bulk_state(DeleteBulkStateRequest {
store_name: store_name.into(),
states: states.into_iter().map(|pair| pair.into()).collect(),
})
.await
}
pub async fn delete_state<S>(
&mut self,
store_name: S,
key: S,
metadata: Option<HashMap<String, String>>,
) -> Result<(), Error>
where
S: Into<String>,
{
let mut mdata = HashMap::<String, String>::new();
if let Some(m) = metadata {
mdata = m;
}
self.0
.delete_state(DeleteStateRequest {
store_name: store_name.into(),
key: key.into(),
metadata: mdata,
..Default::default()
})
.await
}
pub async fn set_metadata<S>(&mut self, key: S, value: S) -> Result<(), Error>
where
S: Into<String>,
{
self.0
.set_metadata(SetMetadataRequest {
key: key.into(),
value: value.into(),
})
.await
}
pub async fn get_metadata(&mut self) -> Result<GetMetadataResponse, Error> {
self.0.get_metadata().await
}
pub async fn invoke_actor<I, M, TInput, TOutput>(
&mut self,
actor_type: I,
actor_id: I,
method_name: M,
input: TInput,
metadata: Option<HashMap<String, String>>,
) -> Result<TOutput, Error>
where
I: Into<String>,
M: Into<String>,
TInput: Serialize,
TOutput: for<'a> Deserialize<'a>,
{
let mut mdata = HashMap::<String, String>::new();
if let Some(m) = metadata {
mdata = m;
}
mdata.insert("Content-Type".to_string(), "application/json".to_string());
let data = match serde_json::to_vec(&input) {
Ok(data) => data,
Err(_e) => return Err(Error::SerializationError),
};
let res = self
.0
.invoke_actor(InvokeActorRequest {
actor_type: actor_type.into(),
actor_id: actor_id.into(),
method: method_name.into(),
data,
metadata: mdata,
})
.await?;
match serde_json::from_slice::<TOutput>(&res.data) {
Ok(output) => Ok(output),
Err(_e) => Err(Error::SerializationError),
}
}
pub async fn get_configuration<S, K>(
&mut self,
store_name: S,
keys: Vec<K>,
metadata: Option<HashMap<String, String>>,
) -> Result<GetConfigurationResponse, Error>
where
S: Into<String>,
K: Into<String>,
{
let request = GetConfigurationRequest {
store_name: store_name.into(),
keys: keys.into_iter().map(|key| key.into()).collect(),
metadata: metadata.unwrap_or_default(),
};
self.0.get_configuration(request).await
}
pub async fn subscribe_configuration<S>(
&mut self,
store_name: S,
keys: Vec<S>,
metadata: Option<HashMap<String, String>>,
) -> Result<Streaming<SubscribeConfigurationResponse>, Error>
where
S: Into<String>,
{
let request = SubscribeConfigurationRequest {
store_name: store_name.into(),
keys: keys.into_iter().map(|key| key.into()).collect(),
metadata: metadata.unwrap_or_default(),
};
self.0.subscribe_configuration(request).await
}
pub async fn unsubscribe_configuration<S>(
&mut self,
store_name: S,
id: S,
) -> Result<UnsubscribeConfigurationResponse, Error>
where
S: Into<String>,
{
let request = UnsubscribeConfigurationRequest {
id: id.into(),
store_name: store_name.into(),
};
self.0.unsubscribe_configuration(request).await
}
pub async fn encrypt<R>(
&mut self,
payload: ReaderStream<R>,
request_options: EncryptRequestOptions,
) -> Result<Vec<StreamPayload>, Status>
where
R: AsyncRead + Send,
{
let request_options = &Some(request_options);
let requested_items: Vec<EncryptRequest> = payload
.0
.enumerate()
.fold(vec![], |mut init, (i, bytes)| async move {
let stream_payload = StreamPayload {
data: bytes.unwrap().to_vec(),
seq: 0,
};
if i == 0 {
init.push(EncryptRequest {
options: request_options.clone(),
payload: Some(stream_payload),
});
} else {
init.push(EncryptRequest {
options: None,
payload: Some(stream_payload),
});
}
init
})
.await;
self.0.encrypt(requested_items).await
}
pub async fn decrypt(
&mut self,
encrypted: Vec<StreamPayload>,
options: DecryptRequestOptions,
) -> Result<Vec<u8>, Status> {
let requested_items: Vec<DecryptRequest> = encrypted
.iter()
.enumerate()
.map(|(i, item)| {
if i == 0 {
DecryptRequest {
options: Some(options.clone()),
payload: Some(item.clone()),
}
} else {
DecryptRequest {
options: None,
payload: Some(item.clone()),
}
}
})
.collect();
self.0.decrypt(requested_items).await
}
pub async fn schedule_job_alpha1(
&mut self,
job: Job,
overwrite: Option<bool>,
) -> Result<ScheduleJobResponse, Error> {
let request = ScheduleJobRequest {
job: Some(job.clone()),
overwrite: overwrite.unwrap_or(false),
};
self.0.schedule_job_alpha1(request).await
}
pub async fn get_job_alpha1(&mut self, name: &str) -> Result<GetJobResponse, Error> {
let request = GetJobRequest {
name: name.to_string(),
};
self.0.get_job_alpha1(request).await
}
pub async fn delete_job_alpha1(&mut self, name: &str) -> Result<DeleteJobResponse, Error> {
let request = DeleteJobRequest {
name: name.to_string(),
};
self.0.delete_job_alpha1(request).await
}
pub async fn converse_alpha1(
&mut self,
request: ConversationRequest,
) -> Result<ConversationResponse, Error> {
self.0.converse_alpha1(request).await
}
pub async fn converse_alpha2(
&mut self,
request: ConversationRequestAlpha2,
) -> Result<ConversationResponseAlpha2, Error> {
self.0.converse_alpha2(request).await
}
}
#[async_trait]
pub trait DaprInterface: Sized {
async fn connect(addr: String) -> Result<Self, Error>;
async fn publish_event(&mut self, request: PublishEventRequest) -> Result<(), Error>;
async fn invoke_service(
&mut self,
request: InvokeServiceRequest,
) -> Result<InvokeServiceResponse, Error>;
async fn invoke_binding(
&mut self,
request: InvokeBindingRequest,
) -> Result<InvokeBindingResponse, Error>;
async fn get_secret(&mut self, request: GetSecretRequest) -> Result<GetSecretResponse, Error>;
async fn get_bulk_secret(
&mut self,
request: GetBulkSecretRequest,
) -> Result<GetBulkSecretResponse, Error>;
async fn get_state(&mut self, request: GetStateRequest) -> Result<GetStateResponse, Error>;
async fn save_state(&mut self, request: SaveStateRequest) -> Result<(), Error>;
async fn query_state_alpha1(
&mut self,
request: QueryStateRequest,
) -> Result<QueryStateResponse, Error>;
async fn delete_state(&mut self, request: DeleteStateRequest) -> Result<(), Error>;
async fn delete_bulk_state(&mut self, request: DeleteBulkStateRequest) -> Result<(), Error>;
async fn set_metadata(&mut self, request: SetMetadataRequest) -> Result<(), Error>;
async fn get_metadata(&mut self) -> Result<GetMetadataResponse, Error>;
async fn invoke_actor(
&mut self,
request: InvokeActorRequest,
) -> Result<InvokeActorResponse, Error>;
async fn get_configuration(
&mut self,
request: GetConfigurationRequest,
) -> Result<GetConfigurationResponse, Error>;
async fn subscribe_configuration(
&mut self,
request: SubscribeConfigurationRequest,
) -> Result<Streaming<SubscribeConfigurationResponse>, Error>;
async fn unsubscribe_configuration(
&mut self,
request: UnsubscribeConfigurationRequest,
) -> Result<UnsubscribeConfigurationResponse, Error>;
async fn encrypt(&mut self, payload: Vec<EncryptRequest>)
-> Result<Vec<StreamPayload>, Status>;
async fn decrypt(&mut self, payload: Vec<DecryptRequest>) -> Result<Vec<u8>, Status>;
async fn schedule_job_alpha1(
&mut self,
request: ScheduleJobRequest,
) -> Result<ScheduleJobResponse, Error>;
async fn get_job_alpha1(&mut self, request: GetJobRequest) -> Result<GetJobResponse, Error>;
async fn delete_job_alpha1(
&mut self,
request: DeleteJobRequest,
) -> Result<DeleteJobResponse, Error>;
async fn converse_alpha1(
&mut self,
request: ConversationRequest,
) -> Result<ConversationResponse, Error>;
async fn converse_alpha2(
&mut self,
request: ConversationRequestAlpha2,
) -> Result<ConversationResponseAlpha2, Error>;
}
async fn connect_plain(
addr: String,
) -> Result<dapr_v1::dapr_client::DaprClient<TonicChannel>, Error> {
Ok(dapr_v1::dapr_client::DaprClient::connect(addr).await?)
}
async fn connect_intercepted(
addr: String,
) -> Result<
dapr_v1::dapr_client::DaprClient<InterceptedService<TonicChannel, ApiTokenInterceptor>>,
Error,
> {
Err(Error::InvalidEndpoint(
crate::error::sanitize_endpoint_for_diagnostics(&addr),
))
}
macro_rules! impl_dapr_interface_for {
($type:ty, $connect_fn:path) => {
#[async_trait]
impl DaprInterface for $type {
async fn connect(addr: String) -> Result<Self, Error> {
$connect_fn(addr).await
}
async fn publish_event(&mut self, request: PublishEventRequest) -> Result<(), Error> {
self.publish_event(Request::new(request))
.await?
.into_inner();
Ok(())
}
async fn invoke_service(
&mut self,
request: InvokeServiceRequest,
) -> Result<InvokeServiceResponse, Error> {
Ok(self
.invoke_service(Request::new(request))
.await?
.into_inner())
}
async fn invoke_binding(
&mut self,
request: InvokeBindingRequest,
) -> Result<InvokeBindingResponse, Error> {
Ok(self
.invoke_binding(Request::new(request))
.await?
.into_inner())
}
async fn get_secret(
&mut self,
request: GetSecretRequest,
) -> Result<GetSecretResponse, Error> {
Ok(self.get_secret(Request::new(request)).await?.into_inner())
}
async fn get_bulk_secret(
&mut self,
request: GetBulkSecretRequest,
) -> Result<GetBulkSecretResponse, Error> {
Ok(self
.get_bulk_secret(Request::new(request))
.await?
.into_inner())
}
async fn get_state(
&mut self,
request: GetStateRequest,
) -> Result<GetStateResponse, Error> {
Ok(self.get_state(Request::new(request)).await?.into_inner())
}
async fn save_state(&mut self, request: SaveStateRequest) -> Result<(), Error> {
self.save_state(Request::new(request)).await?.into_inner();
Ok(())
}
async fn query_state_alpha1(
&mut self,
request: QueryStateRequest,
) -> Result<QueryStateResponse, Error> {
Ok(self
.query_state_alpha1(Request::new(request))
.await?
.into_inner())
}
async fn delete_state(&mut self, request: DeleteStateRequest) -> Result<(), Error> {
self.delete_state(Request::new(request)).await?.into_inner();
Ok(())
}
async fn delete_bulk_state(
&mut self,
request: DeleteBulkStateRequest,
) -> Result<(), Error> {
self.delete_bulk_state(Request::new(request))
.await?
.into_inner();
Ok(())
}
async fn set_metadata(&mut self, request: SetMetadataRequest) -> Result<(), Error> {
self.set_metadata(Request::new(request)).await?.into_inner();
Ok(())
}
async fn get_metadata(&mut self) -> Result<GetMetadataResponse, Error> {
Ok(self.get_metadata(GetMetadataRequest {}).await?.into_inner())
}
async fn invoke_actor(
&mut self,
request: InvokeActorRequest,
) -> Result<InvokeActorResponse, Error> {
Ok(self.invoke_actor(Request::new(request)).await?.into_inner())
}
async fn get_configuration(
&mut self,
request: GetConfigurationRequest,
) -> Result<GetConfigurationResponse, Error> {
Ok(self
.get_configuration(Request::new(request))
.await?
.into_inner())
}
async fn subscribe_configuration(
&mut self,
request: SubscribeConfigurationRequest,
) -> Result<Streaming<SubscribeConfigurationResponse>, Error> {
Ok(self
.subscribe_configuration(Request::new(request))
.await?
.into_inner())
}
async fn unsubscribe_configuration(
&mut self,
request: UnsubscribeConfigurationRequest,
) -> Result<UnsubscribeConfigurationResponse, Error> {
Ok(self
.unsubscribe_configuration(Request::new(request))
.await?
.into_inner())
}
async fn encrypt(
&mut self,
request: Vec<EncryptRequest>,
) -> Result<Vec<StreamPayload>, Status> {
let request = Request::new(tokio_stream::iter(request));
let stream = self.encrypt_alpha1(request).await?;
let mut stream = stream.into_inner();
let mut return_data = vec![];
while let Some(resp) = stream.next().await {
if let Ok(resp) = resp
&& let Some(data) = resp.payload
{
return_data.push(data)
}
}
Ok(return_data)
}
async fn decrypt(&mut self, request: Vec<DecryptRequest>) -> Result<Vec<u8>, Status> {
let request = Request::new(tokio_stream::iter(request));
let stream = self.decrypt_alpha1(request).await?;
let mut stream = stream.into_inner();
let mut data = vec![];
while let Some(resp) = stream.next().await {
if let Ok(resp) = resp
&& let Some(mut payload) = resp.payload
{
data.append(payload.data.as_mut())
}
}
Ok(data)
}
async fn schedule_job_alpha1(
&mut self,
request: ScheduleJobRequest,
) -> Result<ScheduleJobResponse, Error> {
Ok(self.schedule_job_alpha1(request).await?.into_inner())
}
async fn get_job_alpha1(
&mut self,
request: GetJobRequest,
) -> Result<GetJobResponse, Error> {
Ok(self
.get_job_alpha1(Request::new(request))
.await?
.into_inner())
}
async fn delete_job_alpha1(
&mut self,
request: DeleteJobRequest,
) -> Result<DeleteJobResponse, Error> {
Ok(self
.delete_job_alpha1(Request::new(request))
.await?
.into_inner())
}
async fn converse_alpha1(
&mut self,
request: ConversationRequest,
) -> Result<ConversationResponse, Error> {
Ok(self
.converse_alpha1(Request::new(request))
.await?
.into_inner())
}
async fn converse_alpha2(
&mut self,
request: ConversationRequestAlpha2,
) -> Result<ConversationResponseAlpha2, Error> {
Ok(self
.converse_alpha2(Request::new(request))
.await?
.into_inner())
}
}
};
}
impl_dapr_interface_for!(
dapr_v1::dapr_client::DaprClient<TonicChannel>,
connect_plain
);
impl_dapr_interface_for!(
dapr_v1::dapr_client::DaprClient<InterceptedService<TonicChannel, ApiTokenInterceptor>>,
connect_intercepted
);
pub type InvokeServiceRequest = dapr_v1::InvokeServiceRequest;
pub type InvokeServiceResponse = common_v1::InvokeResponse;
pub type InvokeBindingRequest = dapr_v1::InvokeBindingRequest;
pub type InvokeBindingResponse = dapr_v1::InvokeBindingResponse;
pub type PublishEventRequest = dapr_v1::PublishEventRequest;
pub type GetStateRequest = dapr_v1::GetStateRequest;
pub type GetStateResponse = dapr_v1::GetStateResponse;
pub type SaveStateRequest = dapr_v1::SaveStateRequest;
pub type StateItem = common_v1::StateItem;
pub type StateOptions = common_v1::StateOptions;
pub type Etag = common_v1::Etag;
pub type QueryStateRequest = dapr_v1::QueryStateRequest;
pub type QueryStateResponse = dapr_v1::QueryStateResponse;
pub type DeleteStateRequest = dapr_v1::DeleteStateRequest;
pub type DeleteBulkStateRequest = dapr_v1::DeleteBulkStateRequest;
pub type GetSecretRequest = dapr_v1::GetSecretRequest;
pub type GetSecretResponse = dapr_v1::GetSecretResponse;
pub type GetBulkSecretRequest = dapr_v1::GetBulkSecretRequest;
pub type GetBulkSecretResponse = dapr_v1::GetBulkSecretResponse;
pub type GetMetadataResponse = dapr_v1::GetMetadataResponse;
pub type GetMetadataRequest = dapr_v1::GetMetadataRequest;
pub type SetMetadataRequest = dapr_v1::SetMetadataRequest;
pub type InvokeActorRequest = dapr_v1::InvokeActorRequest;
pub type InvokeActorResponse = dapr_v1::InvokeActorResponse;
pub type GetConfigurationRequest = dapr_v1::GetConfigurationRequest;
pub type GetConfigurationResponse = dapr_v1::GetConfigurationResponse;
pub type SubscribeConfigurationRequest = dapr_v1::SubscribeConfigurationRequest;
pub type SubscribeConfigurationResponse = dapr_v1::SubscribeConfigurationResponse;
pub type UnsubscribeConfigurationRequest = dapr_v1::UnsubscribeConfigurationRequest;
pub type TonicClient = dapr_v1::dapr_client::DaprClient<TonicChannel>;
pub type TonicClientWithAuth =
dapr_v1::dapr_client::DaprClient<InterceptedService<TonicChannel, ApiTokenInterceptor>>;
impl Client<TonicClientWithAuth> {
pub async fn new() -> Result<Self, Error> {
let opts = ClientOptions::from_env()?;
Self::from_options(opts).await
}
pub async fn from_options(opts: ClientOptions) -> Result<Self, Error> {
let address = opts.address().to_string();
let interceptor = ApiTokenInterceptor::try_new(opts.api_token().map(|s| s.to_string()))?;
let sanitized_address = crate::error::sanitize_endpoint_for_diagnostics(&address);
let endpoint = tonic::transport::Endpoint::from_shared(address.clone())
.map_err(|_| Error::InvalidEndpoint(sanitized_address.clone()))?
.connect_timeout(opts.timeout());
let channel = match tokio::time::timeout(opts.timeout(), endpoint.connect()).await {
Ok(Ok(c)) => c,
Ok(Err(e)) => return Err(Error::from(e)),
Err(_) => return Err(Error::ConnectTimeout),
};
let grpc = dapr_v1::dapr_client::DaprClient::with_interceptor(channel, interceptor);
Ok(Client(grpc, address))
}
pub async fn connect_with_address(address: impl Into<String>) -> Result<Self, Error> {
let opts = ClientOptions::from_env()?.with_address(address);
Self::from_options(opts).await
}
}
pub type UnsubscribeConfigurationResponse = dapr_v1::UnsubscribeConfigurationResponse;
pub type EncryptRequest = crate::dapr::proto::runtime::v1::EncryptRequest;
pub type DecryptRequest = crate::dapr::proto::runtime::v1::DecryptRequest;
pub type EncryptRequestOptions = crate::dapr::proto::runtime::v1::EncryptRequestOptions;
pub type DecryptRequestOptions = crate::dapr::proto::runtime::v1::DecryptRequestOptions;
pub type Job = crate::dapr::proto::runtime::v1::Job;
pub type JobFailurePolicy = crate::dapr::proto::common::v1::JobFailurePolicy;
pub type ScheduleJobRequest = crate::dapr::proto::runtime::v1::ScheduleJobRequest;
pub type ScheduleJobResponse = crate::dapr::proto::runtime::v1::ScheduleJobResponse;
pub type GetJobRequest = crate::dapr::proto::runtime::v1::GetJobRequest;
pub type GetJobResponse = crate::dapr::proto::runtime::v1::GetJobResponse;
pub type DeleteJobRequest = crate::dapr::proto::runtime::v1::DeleteJobRequest;
pub type DeleteJobResponse = crate::dapr::proto::runtime::v1::DeleteJobResponse;
pub type ConversationRequest = crate::dapr::proto::runtime::v1::ConversationRequest;
pub type ConversationResponse = crate::dapr::proto::runtime::v1::ConversationResponse;
pub type ConversationResult = crate::dapr::proto::runtime::v1::ConversationResult;
pub type ConversationInput = crate::dapr::proto::runtime::v1::ConversationInput;
pub type ConversationRequestAlpha2 = crate::dapr::proto::runtime::v1::ConversationRequestAlpha2;
pub type ConversationResponseAlpha2 = crate::dapr::proto::runtime::v1::ConversationResponseAlpha2;
pub type ConversationInputAlpha2 = crate::dapr::proto::runtime::v1::ConversationInputAlpha2;
type StreamPayload = crate::dapr::proto::common::v1::StreamPayload;
impl<K> From<(K, Vec<u8>)> for common_v1::StateItem
where
K: Into<String>,
{
fn from((key, value): (K, Vec<u8>)) -> Self {
common_v1::StateItem {
key: key.into(),
value,
..Default::default()
}
}
}
pub struct ReaderStream<T>(tokio_util::io::ReaderStream<T>);
impl<T: AsyncRead> ReaderStream<T> {
pub fn new(data: T) -> Self {
ReaderStream(tokio_util::io::ReaderStream::new(data))
}
}
pub struct JobBuilder {
schedule: Option<String>,
data: Option<Any>,
name: String,
ttl: Option<String>,
repeats: Option<u32>,
due_time: Option<String>,
failure_policy: Option<JobFailurePolicy>,
}
impl JobBuilder {
pub fn new(name: &str) -> Self {
JobBuilder {
schedule: None,
data: None,
name: name.to_string(),
ttl: None,
repeats: None,
due_time: None,
failure_policy: None,
}
}
pub fn with_schedule(mut self, schedule: &str) -> Self {
self.schedule = Some(schedule.into());
self
}
pub fn with_data(mut self, data: Any) -> Self {
self.data = Some(data);
self
}
pub fn with_ttl(mut self, ttl: &str) -> Self {
self.ttl = Some(ttl.into());
self
}
pub fn with_repeats(mut self, repeats: u32) -> Self {
self.repeats = Some(repeats);
self
}
pub fn with_due_time(mut self, due_time: &str) -> Self {
self.due_time = Some(due_time.into());
self
}
pub fn with_failure_policy(mut self, policy: JobFailurePolicy) -> Self {
self.failure_policy = Some(policy);
self
}
pub fn build(self) -> Job {
Job {
schedule: self.schedule,
data: self.data,
name: self.name,
ttl: self.ttl,
repeats: self.repeats,
due_time: self.due_time,
failure_policy: self.failure_policy,
}
}
}
pub enum JobFailurePolicyType {
Drop {},
Constant {},
}
pub struct JobFailurePolicyBuilder {
policy: JobFailurePolicyType,
pub retry_interval: Option<Duration>,
pub max_retries: Option<u32>,
}
impl JobFailurePolicyBuilder {
pub fn new(policy: JobFailurePolicyType) -> Self {
JobFailurePolicyBuilder {
policy,
retry_interval: None,
max_retries: None,
}
}
pub fn with_retry_interval(mut self, interval: Duration) -> Self {
self.retry_interval = Some(interval);
self
}
pub fn with_max_retries(mut self, max_retries: u32) -> Self {
self.max_retries = Some(max_retries);
self
}
pub fn build(self) -> common_v1::JobFailurePolicy {
match self.policy {
JobFailurePolicyType::Drop {} => common_v1::JobFailurePolicy {
policy: Some(Policy::Drop(Default::default())),
},
JobFailurePolicyType::Constant {} => JobFailurePolicy {
policy: Some(Policy::Constant(JobFailurePolicyConstant {
interval: self.retry_interval.map(|interval| {
prost_types::Duration::try_from(interval)
.expect("Failed to convert Duration")
}),
max_retries: self.max_retries,
})),
},
}
}
}
pub struct ConversationInputBuilder {
content: String,
role: Option<String>,
scrub_pii: Option<bool>,
}
impl ConversationInputBuilder {
pub fn new(message: &str) -> Self {
ConversationInputBuilder {
content: message.to_string(),
role: None,
scrub_pii: None,
}
}
pub fn build(self) -> ConversationInput {
ConversationInput {
content: self.content,
role: self.role,
scrub_pii: self.scrub_pii,
}
}
}
pub struct ConversationRequestBuilder {
name: String,
context_id: Option<String>,
inputs: Vec<ConversationInput>,
parameters: HashMap<String, Any>,
metadata: HashMap<String, String>,
scrub_pii: Option<bool>,
temperature: Option<f64>,
}
impl ConversationRequestBuilder {
pub fn new(name: &str, inputs: Vec<ConversationInput>) -> Self {
ConversationRequestBuilder {
name: name.to_string(),
context_id: None,
inputs,
parameters: Default::default(),
metadata: Default::default(),
scrub_pii: None,
temperature: None,
}
}
pub fn build(self) -> ConversationRequest {
ConversationRequest {
name: self.name,
context_id: self.context_id,
inputs: self.inputs,
parameters: self.parameters,
metadata: self.metadata,
scrub_pii: self.scrub_pii,
temperature: self.temperature,
}
}
}
pub struct ConversationInputAlpha2Builder {
messages: Vec<crate::dapr::proto::runtime::v1::ConversationMessage>,
scrub_pii: Option<bool>,
}
impl ConversationInputAlpha2Builder {
pub fn new(messages: Vec<crate::dapr::proto::runtime::v1::ConversationMessage>) -> Self {
ConversationInputAlpha2Builder {
messages,
scrub_pii: None,
}
}
pub fn with_scrub_pii(mut self, scrub: bool) -> Self {
self.scrub_pii = Some(scrub);
self
}
pub fn build(self) -> ConversationInputAlpha2 {
ConversationInputAlpha2 {
messages: self.messages,
scrub_pii: self.scrub_pii,
}
}
}
pub struct ConversationRequestAlpha2Builder {
name: String,
inputs: Vec<ConversationInputAlpha2>,
metadata: HashMap<String, String>,
scrub_pii: Option<bool>,
temperature: Option<f64>,
parameters: HashMap<String, prost_types::Any>,
tools: Vec<crate::dapr::proto::runtime::v1::ConversationTools>,
response_format: Option<prost_types::Struct>,
prompt_cache_retention: Option<prost_types::Duration>,
}
impl ConversationRequestAlpha2Builder {
pub fn new(name: &str, inputs: Vec<ConversationInputAlpha2>) -> Self {
ConversationRequestAlpha2Builder {
name: name.to_string(),
inputs,
metadata: Default::default(),
scrub_pii: None,
temperature: None,
parameters: Default::default(),
tools: Default::default(),
response_format: None,
prompt_cache_retention: None,
}
}
pub fn with_metadata(mut self, metadata: HashMap<String, String>) -> Self {
self.metadata = metadata;
self
}
pub fn with_scrub_pii(mut self, scrub: bool) -> Self {
self.scrub_pii = Some(scrub);
self
}
pub fn with_temperature(mut self, temp: f64) -> Self {
self.temperature = Some(temp);
self
}
pub fn with_parameters(mut self, params: HashMap<String, prost_types::Any>) -> Self {
self.parameters = params;
self
}
pub fn with_tools(
mut self,
tools: Vec<crate::dapr::proto::runtime::v1::ConversationTools>,
) -> Self {
self.tools = tools;
self
}
pub fn with_response_format(mut self, format: prost_types::Struct) -> Self {
self.response_format = Some(format);
self
}
pub fn with_prompt_cache_retention(mut self, retention: prost_types::Duration) -> Self {
self.prompt_cache_retention = Some(retention);
self
}
pub fn build(self) -> ConversationRequestAlpha2 {
ConversationRequestAlpha2 {
name: self.name,
inputs: self.inputs,
metadata: self.metadata,
scrub_pii: self.scrub_pii,
temperature: self.temperature,
parameters: self.parameters,
tools: self.tools,
response_format: self.response_format,
prompt_cache_retention: self.prompt_cache_retention,
..Default::default()
}
}
}
pub use crate::dapr::proto::runtime::v1::{
ConversationMessage, ConversationMessageContent, ConversationMessageOfAssistant,
ConversationMessageOfDeveloper, ConversationMessageOfSystem, ConversationMessageOfTool,
ConversationMessageOfUser, ConversationToolCalls, ConversationToolCallsOfFunction,
ConversationTools, conversation_message, conversation_tool_calls,
};
impl From<ConversationMessageOfUser> for ConversationMessage {
fn from(msg: ConversationMessageOfUser) -> Self {
ConversationMessage {
message_types: Some(conversation_message::MessageTypes::OfUser(msg)),
}
}
}
impl From<ConversationMessageOfSystem> for ConversationMessage {
fn from(msg: ConversationMessageOfSystem) -> Self {
ConversationMessage {
message_types: Some(conversation_message::MessageTypes::OfSystem(msg)),
}
}
}
impl From<ConversationMessageOfDeveloper> for ConversationMessage {
fn from(msg: ConversationMessageOfDeveloper) -> Self {
ConversationMessage {
message_types: Some(conversation_message::MessageTypes::OfDeveloper(msg)),
}
}
}
impl From<ConversationMessageOfAssistant> for ConversationMessage {
fn from(msg: ConversationMessageOfAssistant) -> Self {
ConversationMessage {
message_types: Some(conversation_message::MessageTypes::OfAssistant(msg)),
}
}
}
impl From<ConversationMessageOfTool> for ConversationMessage {
fn from(msg: ConversationMessageOfTool) -> Self {
ConversationMessage {
message_types: Some(conversation_message::MessageTypes::OfTool(msg)),
}
}
}