use crate::{
ClientInterceptor, HttpConnectProxyOptions, RetryOptions, RpcOptions, VERSION, callback_based,
};
#[cfg(feature = "experimental")]
use crate::{ClientPlugin, ErasedClientPlugin};
use http::Uri;
use std::{collections::HashMap, sync::Arc, time::Duration};
use temporalio_common::{
ActivityCloseTimeouts, MemoValues, RetryPolicy,
data_converters::{
DataConverter, GenericPayloadConverter, PayloadConversionError, PayloadConverter,
SerializationContext, SerializationContextData, WorkflowSerializationContext,
},
payload_visitor::encode_payloads,
protos::temporal::api::{
common::{
self,
v1::{Header, Memo as ProtoMemo, Payloads},
},
enums::v1::{
ActivityIdConflictPolicy as ProtoActivityIdConflictPolicy,
ActivityIdReusePolicy as ProtoActivityIdReusePolicy,
ArchivalState as ProtoArchivalState,
HistoryEventFilterType as ProtoHistoryEventFilterType,
QueryRejectCondition as ProtoQueryRejectCondition,
WorkflowIdConflictPolicy as ProtoWorkflowIdConflictPolicy,
WorkflowIdReusePolicy as ProtoWorkflowIdReusePolicy,
},
replication::v1::ClusterReplicationConfig,
sdk::v1::UserMetadata,
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;
pub(crate) const DEFAULT_PAYLOADS_WARN_SIZE: u64 = 512 * 1024;
pub(crate) const DEFAULT_MEMO_WARN_SIZE: u64 = 2 * 1024;
#[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,
#[cfg(feature = "experimental")]
#[cfg_attr(
docsrs,
builder(setters(
some_fn(name = payload_limits_impl, vis = "pub(crate)"),
option_fn(name = maybe_payload_limits_impl, vis = "pub(crate)")
))
)]
#[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(all(feature = "experimental", docsrs))]
impl<S: connection_options_builder::State> ConnectionOptionsBuilder<S> {
#[doc(cfg(feature = "experimental"))]
pub fn payload_limits(
self,
value: PayloadLimitsOptions,
) -> ConnectionOptionsBuilder<connection_options_builder::SetPayloadLimits<S>>
where
S::PayloadLimits: connection_options_builder::IsUnset,
{
self.payload_limits_impl(value)
}
#[doc(cfg(feature = "experimental"))]
pub fn maybe_payload_limits(
self,
value: Option<PayloadLimitsOptions>,
) -> ConnectionOptionsBuilder<connection_options_builder::SetPayloadLimits<S>>
where
S::PayloadLimits: connection_options_builder::IsUnset,
{
self.maybe_payload_limits_impl(value)
}
}
#[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)]
#[cfg(feature = "experimental")]
plugins: Vec<ErasedClientPlugin>,
#[builder(field)]
#[debug(skip)]
#[cfg(feature = "experimental")]
client_plugins_applied: bool,
#[builder(default)]
pub data_converter: DataConverter,
#[builder(default)]
#[debug(skip)]
pub client_interceptors: Vec<Arc<dyn ClientInterceptor>>,
}
#[cfg(feature = "experimental")]
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 {
#[cfg(feature = "experimental")]
pub fn plugins(&self) -> &[ErasedClientPlugin] {
&self.plugins
}
#[cfg(feature = "experimental")]
pub(crate) fn client_plugins_applied(&self) -> bool {
self.client_plugins_applied
}
#[cfg(feature = "experimental")]
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()
}
}
#[cfg(feature = "experimental")]
#[derive(Clone, Debug, PartialEq, bon::Builder)]
#[non_exhaustive]
pub struct PayloadLimitsOptions {
#[builder(default = DEFAULT_PAYLOADS_WARN_SIZE)]
pub payloads_warn_size: u64,
#[builder(default = DEFAULT_MEMO_WARN_SIZE)]
pub memo_warn_size: u64,
}
#[cfg(feature = "experimental")]
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, Copy, PartialEq, Eq, Hash, Default)]
#[non_exhaustive]
pub enum WorkflowIdReusePolicy {
#[default]
Unspecified,
AllowDuplicate,
AllowDuplicateFailedOnly,
RejectDuplicate,
}
impl From<WorkflowIdReusePolicy> for ProtoWorkflowIdReusePolicy {
fn from(value: WorkflowIdReusePolicy) -> Self {
match value {
WorkflowIdReusePolicy::Unspecified => Self::Unspecified,
WorkflowIdReusePolicy::AllowDuplicate => Self::AllowDuplicate,
WorkflowIdReusePolicy::AllowDuplicateFailedOnly => Self::AllowDuplicateFailedOnly,
WorkflowIdReusePolicy::RejectDuplicate => Self::RejectDuplicate,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
#[non_exhaustive]
pub enum WorkflowIdConflictPolicy {
#[default]
Unspecified,
Fail,
UseExisting,
TerminateExisting,
}
impl From<WorkflowIdConflictPolicy> for ProtoWorkflowIdConflictPolicy {
fn from(value: WorkflowIdConflictPolicy) -> Self {
match value {
WorkflowIdConflictPolicy::Unspecified => Self::Unspecified,
WorkflowIdConflictPolicy::Fail => Self::Fail,
WorkflowIdConflictPolicy::UseExisting => Self::UseExisting,
WorkflowIdConflictPolicy::TerminateExisting => Self::TerminateExisting,
}
}
}
#[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>,
#[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 memo: Option<MemoValues>,
pub static_summary: Option<String>,
pub static_details: Option<String>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
impl WorkflowStartOptions {
pub(crate) async fn encoded_memo(
&self,
data_converter: &DataConverter,
) -> Result<Option<ProtoMemo>, PayloadConversionError> {
let Some(memo) = &self.memo else {
return Ok(None);
};
let payload_converter = data_converter.payload_converter();
let context_data = SerializationContextData::Workflow(WorkflowSerializationContext::new());
let context = SerializationContext::new(&context_data, payload_converter);
let mut memo = ProtoMemo {
fields: memo
.iter()
.map(|(key, value)| {
payload_converter
.to_payload(&context, value)
.map(|payload| (key.to_owned(), payload))
})
.collect::<Result<_, _>>()?,
};
encode_payloads(
&mut memo,
data_converter.codec(),
&SerializationContextData::Workflow(WorkflowSerializationContext::new()),
)
.await?;
Ok(Some(memo))
}
pub(crate) fn user_metadata(&self) -> Option<UserMetadata> {
(self.static_summary.is_some() || self.static_details.is_some()).then(|| {
let payload_converter = PayloadConverter::default();
let context_data =
SerializationContextData::Workflow(WorkflowSerializationContext::new());
let context = SerializationContext::new(&context_data, &payload_converter);
UserMetadata {
summary: self.static_summary.as_ref().map(|summary| {
payload_converter
.to_payload(&context, summary)
.expect("String-to-JSON payload serialization is infallible")
}),
details: self.static_details.as_ref().map(|details| {
payload_converter
.to_payload(&context, details)
.expect("String-to-JSON payload serialization is infallible")
}),
}
})
}
}
#[derive(Debug, Clone, bon::Builder)]
#[builder(start_fn = new, on(String, into))]
#[non_exhaustive]
pub struct WorkflowUpdateWithStartOptions {
#[builder(start_fn)]
pub task_queue: String,
#[builder(start_fn)]
pub workflow_id: String,
#[builder(start_fn)]
pub id_conflict_policy: WorkflowIdConflictPolicy,
#[builder(default)]
pub id_reuse_policy: WorkflowIdReusePolicy,
pub execution_timeout: Option<Duration>,
pub run_timeout: Option<Duration>,
pub task_timeout: Option<Duration>,
pub search_attributes: Option<SearchAttributes>,
#[builder(into)]
pub retry_policy: Option<RetryPolicy>,
#[builder(default)]
pub links: Vec<common::v1::Link>,
#[builder(default)]
pub completion_callbacks: Vec<common::v1::Callback>,
#[builder(default)]
pub priority: Priority,
pub start_header: Option<Header>,
pub update_header: Option<Header>,
pub memo: Option<MemoValues>,
pub static_summary: Option<String>,
pub static_details: Option<String>,
pub update_id: Option<String>,
#[builder(default)]
pub rpc_options: RpcOptions,
}
impl WorkflowUpdateWithStartOptions {
pub(crate) fn into_parts(self) -> (WorkflowStartOptions, Option<String>, Option<Header>) {
let Self {
task_queue,
workflow_id,
id_conflict_policy,
id_reuse_policy,
execution_timeout,
run_timeout,
task_timeout,
search_attributes,
retry_policy,
links,
completion_callbacks,
priority,
start_header,
update_header,
memo,
static_summary,
static_details,
update_id,
rpc_options: _,
} = self;
(
WorkflowStartOptions {
task_queue,
workflow_id,
id_reuse_policy,
id_conflict_policy,
execution_timeout,
run_timeout,
task_timeout,
cron_schedule: None,
search_attributes,
enable_eager_workflow_start: false,
retry_policy,
links,
completion_callbacks,
priority,
header: start_header,
memo,
static_summary,
static_details,
rpc_options: RpcOptions::default(),
},
update_id,
update_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, Copy, PartialEq, Eq, Hash, Default)]
#[non_exhaustive]
pub enum QueryRejectCondition {
#[default]
Unspecified,
None,
NotOpen,
NotCompletedCleanly,
}
impl From<QueryRejectCondition> for ProtoQueryRejectCondition {
fn from(value: QueryRejectCondition) -> Self {
match value {
QueryRejectCondition::Unspecified => Self::Unspecified,
QueryRejectCondition::None => Self::None,
QueryRejectCondition::NotOpen => Self::NotOpen,
QueryRejectCondition::NotCompletedCleanly => Self::NotCompletedCleanly,
}
}
}
#[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(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
#[non_exhaustive]
pub enum ArchivalState {
#[default]
Unspecified,
Disabled,
Enabled,
}
impl From<ArchivalState> for ProtoArchivalState {
fn from(value: ArchivalState) -> Self {
match value {
ArchivalState::Unspecified => Self::Unspecified,
ArchivalState::Disabled => Self::Disabled,
ArchivalState::Enabled => Self::Enabled,
}
}
}
#[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: ProtoArchivalState::from(val.history_archival_state) as i32,
history_archival_uri: val.history_archival_uri,
visibility_archival_state: ProtoArchivalState::from(val.visibility_archival_state)
as i32,
visibility_archival_uri: val.visibility_archival_uri,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
#[non_exhaustive]
pub enum HistoryEventFilterType {
#[default]
Unspecified,
AllEvent,
CloseEvent,
}
impl From<HistoryEventFilterType> for ProtoHistoryEventFilterType {
fn from(value: HistoryEventFilterType) -> Self {
match value {
HistoryEventFilterType::Unspecified => Self::Unspecified,
HistoryEventFilterType::AllEvent => Self::AllEvent,
HistoryEventFilterType::CloseEvent => Self::CloseEvent,
}
}
}
#[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,
}
impl From<WorkflowExecuteUpdateOptions> for WorkflowStartUpdateOptions {
fn from(options: WorkflowExecuteUpdateOptions) -> Self {
Self::builder()
.maybe_update_id(options.update_id)
.maybe_header(options.header)
.rpc_options(options.rpc_options)
.build()
}
}
#[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,
}