use crate::{
ClientInterceptor, ClientPlugin, ErasedClientPlugin, HttpConnectProxyOptions, RetryOptions,
RpcOptions, VERSION, callback_based,
};
use http::Uri;
use std::{collections::HashMap, sync::Arc, time::Duration};
use temporalio_common::{
ActivityCloseTimeouts, RetryPolicy,
data_converters::DataConverter,
protos::temporal::api::{
common::{
self,
v1::{Header, Payloads},
},
enums::v1::{
ActivityIdConflictPolicy as ProtoActivityIdConflictPolicy,
ActivityIdReusePolicy as ProtoActivityIdReusePolicy, ArchivalState,
HistoryEventFilterType, QueryRejectCondition, WorkflowIdConflictPolicy,
WorkflowIdReusePolicy,
},
replication::v1::ClusterReplicationConfig,
workflowservice::v1::RegisterNamespaceRequest,
},
search_attributes::SearchAttributes,
telemetry::metrics::TemporalMeter,
};
#[cfg(feature = "dynamic-tls")]
use tokio_rustls::rustls::client::ResolvesClientCert;
use tokio_rustls::rustls::client::danger::ServerCertVerifier;
use url::Url;
#[derive(bon::Builder, Clone, Debug)]
#[non_exhaustive]
#[builder(start_fn = new, on(String, into), state_mod(vis = "pub"))]
pub struct ConnectionOptions {
#[builder(start_fn, into)]
pub target: Url,
#[builder(default)]
pub identity: String,
pub metrics_meter: Option<TemporalMeter>,
pub tls_options: Option<TlsOptions>,
pub override_origin: Option<Uri>,
pub api_key: Option<String>,
pub connect_timeout: Option<Duration>,
#[builder(default)]
pub retry_options: RetryOptions,
#[builder(required, default = Some(ClientKeepAliveOptions::default()))]
pub keep_alive: Option<ClientKeepAliveOptions>,
pub headers: Option<HashMap<String, String>>,
pub binary_headers: Option<HashMap<String, Vec<u8>>>,
pub http_connect_proxy: Option<HttpConnectProxyOptions>,
#[builder(required, default = Some(DnsLoadBalancingOptions::default()))]
pub dns_load_balancing: Option<DnsLoadBalancingOptions>,
#[builder(default)]
pub disable_error_code_metric_tags: bool,
pub service_override: Option<callback_based::CallbackBasedGrpcService>,
#[builder(default)]
pub grpc_compression: GrpcCompression,
#[builder(default)]
pub payload_limits: PayloadLimitsOptions,
#[builder(default)]
#[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
pub(crate) skip_get_system_info: bool,
#[builder(default = "temporal-rust".to_owned())]
#[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
pub(crate) client_name: String,
#[builder(default = VERSION.to_owned())]
#[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
pub(crate) client_version: String,
}
#[cfg(feature = "core-based-sdk")]
impl ConnectionOptions {
pub fn set_skip_get_system_info(&mut self, skip: bool) {
self.skip_get_system_info = skip;
}
pub fn get_skip_get_system_info(&self) -> bool {
self.skip_get_system_info
}
pub fn get_client_name(&self) -> &str {
&self.client_name
}
pub fn get_client_version(&self) -> &str {
&self.client_version
}
}
#[derive(Clone, derive_more::Debug, bon::Builder)]
#[non_exhaustive]
#[builder(start_fn = new, on(String, into), state_mod(vis = "pub"))]
pub struct ClientOptions {
#[builder(start_fn)]
pub namespace: String,
#[builder(field)]
#[debug(skip)]
plugins: Vec<ErasedClientPlugin>,
#[builder(field)]
#[debug(skip)]
client_plugins_applied: bool,
#[builder(default)]
pub data_converter: DataConverter,
#[builder(default)]
#[debug(skip)]
pub client_interceptors: Vec<Arc<dyn ClientInterceptor>>,
}
impl<S: client_options_builder::State> ClientOptionsBuilder<S> {
pub fn plugin<P: Into<ErasedClientPlugin>>(mut self, plugin: P) -> Self {
self.plugins.push(plugin.into());
self
}
pub fn plugins<I, P>(mut self, plugins: I) -> Self
where
I: IntoIterator<Item = P>,
P: Into<ErasedClientPlugin>,
{
self.plugins.extend(plugins.into_iter().map(Into::into));
self
}
pub fn client_plugin<P: ClientPlugin>(mut self, plugin: P) -> Self {
self.plugins.push(ErasedClientPlugin::new(plugin));
self
}
}
impl ClientOptions {
pub fn plugins(&self) -> &[ErasedClientPlugin] {
&self.plugins
}
pub(crate) fn client_plugins_applied(&self) -> bool {
self.client_plugins_applied
}
pub(crate) fn mark_client_plugins_applied(&mut self) {
self.client_plugins_applied = true;
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
#[non_exhaustive]
pub enum GrpcCompression {
None,
#[default]
Gzip,
}
#[derive(Clone, bon::Builder)]
#[non_exhaustive]
pub struct TlsOptions {
pub server_root_ca_cert: Option<Vec<u8>>,
pub domain: Option<String>,
pub client_tls_options: Option<ClientTlsOptions>,
pub server_cert_verifier: Option<Arc<dyn ServerCertVerifier>>,
#[cfg(feature = "dynamic-tls")]
pub client_cert_resolver: Option<Arc<dyn ResolvesClientCert>>,
}
impl Default for TlsOptions {
fn default() -> Self {
Self::builder().build()
}
}
impl std::fmt::Debug for TlsOptions {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let mut s = f.debug_struct("TlsOptions");
s.field(
"server_root_ca_cert",
&self
.server_root_ca_cert
.as_ref()
.map(|c| format!("{} bytes", c.len())),
);
s.field("domain", &self.domain);
s.field("client_tls_options", &self.client_tls_options);
s.field(
"server_cert_verifier",
&self.server_cert_verifier.as_ref().map(|_| "<custom>"),
);
#[cfg(feature = "dynamic-tls")]
s.field(
"client_cert_resolver",
&self.client_cert_resolver.as_ref().map(|_| "<custom>"),
);
s.finish()
}
}
#[derive(Clone, bon::Builder)]
#[non_exhaustive]
pub struct ClientTlsOptions {
pub client_cert: Vec<u8>,
pub client_private_key: Vec<u8>,
}
#[derive(Clone, Debug, PartialEq, bon::Builder)]
#[non_exhaustive]
pub struct ClientKeepAliveOptions {
#[builder(default = Duration::from_secs(30))]
pub interval: Duration,
#[builder(default = Duration::from_secs(15))]
pub timeout: Duration,
}
impl Default for ClientKeepAliveOptions {
fn default() -> Self {
Self::builder().build()
}
}
#[derive(Clone, Debug, PartialEq, bon::Builder)]
#[non_exhaustive]
pub struct DnsLoadBalancingOptions {
#[builder(default = Duration::from_secs(30))]
pub resolution_interval: Duration,
}
impl Default for DnsLoadBalancingOptions {
fn default() -> Self {
Self::builder().build()
}
}
#[derive(Clone, Debug, PartialEq, bon::Builder)]
#[non_exhaustive]
pub struct PayloadLimitsOptions {
#[builder(default = 512 * 1024)]
pub payloads_warn_size: u64,
#[builder(default = 2 * 1024)]
pub memo_warn_size: u64,
}
impl Default for PayloadLimitsOptions {
fn default() -> Self {
Self::builder().build()
}
}
impl std::fmt::Debug for ClientTlsOptions {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "ClientTlsOptions(..)")
}
}
#[derive(Debug, Clone, bon::Builder)]
#[builder(start_fn = new, on(String, into))]
#[non_exhaustive]
pub struct WorkflowStartOptions {
#[builder(start_fn)]
pub task_queue: String,
#[builder(start_fn)]
pub workflow_id: String,
#[builder(default)]
pub id_reuse_policy: WorkflowIdReusePolicy,
#[builder(default)]
pub id_conflict_policy: WorkflowIdConflictPolicy,
pub execution_timeout: Option<Duration>,
pub run_timeout: Option<Duration>,
pub task_timeout: Option<Duration>,
pub cron_schedule: Option<String>,
pub search_attributes: Option<SearchAttributes>,
#[builder(default)]
pub enable_eager_workflow_start: bool,
#[builder(into)]
pub retry_policy: Option<RetryPolicy>,
pub start_signal: Option<WorkflowStartSignal>,
#[builder(default)]
pub links: Vec<common::v1::Link>,
#[builder(default)]
pub completion_callbacks: Vec<common::v1::Callback>,
#[builder(default)]
pub priority: Priority,
pub header: Option<Header>,
pub static_summary: Option<String>,
pub static_details: Option<String>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, bon::Builder)]
#[builder(start_fn = new, on(String, into))]
#[non_exhaustive]
pub struct WorkflowStartSignal {
#[builder(start_fn)]
pub signal_name: String,
pub input: Option<Payloads>,
pub header: Option<Header>,
}
pub use temporalio_common::Priority;
#[derive(Debug, Clone, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowGetResultOptions {
#[builder(default = true)]
pub follow_runs: bool,
#[builder(default)]
pub rpc_options: RpcOptions,
}
impl Default for WorkflowGetResultOptions {
fn default() -> Self {
Self {
follow_runs: true,
rpc_options: RpcOptions::default(),
}
}
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowExecuteUpdateOptions {
pub update_id: Option<String>,
pub header: Option<Header>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowSignalOptions {
pub request_id: Option<String>,
pub header: Option<Header>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowQueryOptions {
pub reject_condition: Option<QueryRejectCondition>,
pub header: Option<Header>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[builder(on(String, into))]
#[non_exhaustive]
pub struct WorkflowCancelOptions {
#[builder(default)]
pub reason: String,
pub request_id: Option<String>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[builder(on(String, into))]
#[non_exhaustive]
pub struct WorkflowTerminateOptions {
#[builder(default)]
pub reason: String,
pub details: Option<Payloads>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowDescribeOptions {
#[builder(default)]
pub rpc_options: RpcOptions,
}
const DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD: Duration = Duration::from_secs(60 * 60 * 24 * 3);
#[derive(Clone, Debug, bon::Builder)]
#[builder(on(String, into))]
#[non_exhaustive]
pub struct RegisterNamespaceOptions {
pub namespace: String,
pub description: String,
#[builder(default)]
pub owner_email: String,
#[builder(default = DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD)]
pub workflow_execution_retention_period: Duration,
#[builder(default)]
pub clusters: Vec<ClusterReplicationConfig>,
#[builder(default)]
pub active_cluster_name: String,
#[builder(default)]
pub data: HashMap<String, String>,
#[builder(default)]
pub security_token: String,
#[builder(default)]
pub is_global_namespace: bool,
#[builder(default = ArchivalState::Unspecified)]
pub history_archival_state: ArchivalState,
#[builder(default)]
pub history_archival_uri: String,
#[builder(default = ArchivalState::Unspecified)]
pub visibility_archival_state: ArchivalState,
#[builder(default)]
pub visibility_archival_uri: String,
}
impl From<RegisterNamespaceOptions> for RegisterNamespaceRequest {
fn from(val: RegisterNamespaceOptions) -> Self {
RegisterNamespaceRequest {
namespace: val.namespace,
description: val.description,
owner_email: val.owner_email,
workflow_execution_retention_period: val
.workflow_execution_retention_period
.try_into()
.ok(),
clusters: val.clusters,
active_cluster_name: val.active_cluster_name,
data: val.data,
security_token: val.security_token,
is_global_namespace: val.is_global_namespace,
history_archival_state: val.history_archival_state as i32,
history_archival_uri: val.history_archival_uri,
visibility_archival_state: val.visibility_archival_state as i32,
visibility_archival_uri: val.visibility_archival_uri,
}
}
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowFetchHistoryOptions {
#[builder(default)]
pub skip_archival: bool,
#[builder(default)]
pub wait_new_event: bool,
#[builder(default = HistoryEventFilterType::AllEvent)]
pub event_filter_type: HistoryEventFilterType,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowStartUpdateOptions {
pub update_id: Option<String>,
pub header: Option<Header>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowListOptions {
pub limit: Option<usize>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct WorkflowCountOptions {
#[builder(default)]
pub rpc_options: RpcOptions,
}
#[derive(Clone, Debug, bon::Builder)]
#[builder(start_fn = new, on(String, into))]
#[non_exhaustive]
pub struct ActivityStartOptions {
#[builder(start_fn)]
pub task_queue: String,
#[builder(start_fn)]
pub id: String,
#[builder(start_fn)]
pub close_timeouts: ActivityCloseTimeouts,
pub schedule_to_start_timeout: Option<Duration>,
pub heartbeat_timeout: Option<Duration>,
#[builder(into)]
pub retry_policy: Option<RetryPolicy>,
#[builder(default)]
pub priority: Priority,
#[builder(default)]
pub id_reuse_policy: ActivityIdReusePolicy,
#[builder(default)]
pub id_conflict_policy: ActivityIdConflictPolicy,
pub search_attributes: Option<SearchAttributes>,
pub header: Option<Header>,
pub summary: Option<String>,
pub static_details: Option<String>,
pub start_delay: Option<Duration>,
}
impl ActivityStartOptions {
pub fn with_start_to_close_timeout(
task_queue: impl Into<String>,
activity_id: impl Into<String>,
start_to_close_timeout: Duration,
) -> ActivityStartOptionsBuilder {
Self::new(
task_queue,
activity_id,
ActivityCloseTimeouts::StartToClose(start_to_close_timeout),
)
}
pub fn with_schedule_to_close_timeout(
task_queue: impl Into<String>,
activity_id: impl Into<String>,
schedule_to_close_timeout: Duration,
) -> ActivityStartOptionsBuilder {
Self::new(
task_queue,
activity_id,
ActivityCloseTimeouts::ScheduleToClose(schedule_to_close_timeout),
)
}
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum ActivityIdReusePolicy {
#[default]
AllowDuplicate,
AllowDuplicateFailedOnly,
RejectDuplicate,
}
impl From<ActivityIdReusePolicy> for ProtoActivityIdReusePolicy {
fn from(value: ActivityIdReusePolicy) -> Self {
match value {
ActivityIdReusePolicy::AllowDuplicate => Self::AllowDuplicate,
ActivityIdReusePolicy::AllowDuplicateFailedOnly => Self::AllowDuplicateFailedOnly,
ActivityIdReusePolicy::RejectDuplicate => Self::RejectDuplicate,
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
pub enum ActivityIdConflictPolicy {
#[default]
Fail,
UseExisting,
}
impl From<ActivityIdConflictPolicy> for ProtoActivityIdConflictPolicy {
fn from(value: ActivityIdConflictPolicy) -> Self {
match value {
ActivityIdConflictPolicy::Fail => Self::Fail,
ActivityIdConflictPolicy::UseExisting => Self::UseExisting,
}
}
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct ActivityListOptions {}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct ActivityCountOptions {}
#[derive(Debug, Clone, Default, bon::Builder)]
#[non_exhaustive]
pub struct ActivityDescribeOptions {
#[builder(default)]
pub include_input: bool,
#[builder(default)]
pub include_outcome: bool,
#[builder(default)]
pub include_heartbeat_details: bool,
#[builder(default)]
pub include_last_failure: bool,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[builder(on(String, into))]
#[non_exhaustive]
pub struct ActivityCancelOptions {
#[builder(default)]
pub reason: String,
}
#[derive(Debug, Clone, Default, bon::Builder)]
#[builder(on(String, into))]
#[non_exhaustive]
pub struct ActivityTerminateOptions {
#[builder(default)]
pub reason: String,
}