1use crate::{
2 ClientInterceptor, HttpConnectProxyOptions, RetryOptions, RpcOptions, VERSION, callback_based,
3};
4#[cfg(feature = "experimental")]
5use crate::{ClientPlugin, ErasedClientPlugin};
6use http::Uri;
7use std::{collections::HashMap, sync::Arc, time::Duration};
8use temporalio_common::{
9 ActivityCloseTimeouts, MemoValues, RetryPolicy,
10 data_converters::{
11 DataConverter, GenericPayloadConverter, PayloadConversionError, PayloadConverter,
12 SerializationContext, SerializationContextData, WorkflowSerializationContext,
13 },
14 payload_visitor::encode_payloads,
15 protos::temporal::api::{
16 common::{
17 self,
18 v1::{Header, Memo as ProtoMemo, Payloads},
19 },
20 enums::v1::{
21 ActivityIdConflictPolicy as ProtoActivityIdConflictPolicy,
22 ActivityIdReusePolicy as ProtoActivityIdReusePolicy,
23 ArchivalState as ProtoArchivalState,
24 HistoryEventFilterType as ProtoHistoryEventFilterType,
25 QueryRejectCondition as ProtoQueryRejectCondition,
26 WorkflowIdConflictPolicy as ProtoWorkflowIdConflictPolicy,
27 WorkflowIdReusePolicy as ProtoWorkflowIdReusePolicy,
28 },
29 replication::v1::ClusterReplicationConfig,
30 sdk::v1::UserMetadata,
31 workflowservice::v1::RegisterNamespaceRequest,
32 },
33 search_attributes::SearchAttributes,
34 telemetry::metrics::TemporalMeter,
35};
36#[cfg(feature = "dynamic-tls")]
37use tokio_rustls::rustls::client::ResolvesClientCert;
38use tokio_rustls::rustls::client::danger::ServerCertVerifier;
39use url::Url;
40
41pub(crate) const DEFAULT_PAYLOADS_WARN_SIZE: u64 = 512 * 1024;
42pub(crate) const DEFAULT_MEMO_WARN_SIZE: u64 = 2 * 1024;
43
44#[derive(bon::Builder, Clone, Debug)]
46#[non_exhaustive]
47#[builder(start_fn = new, on(String, into), state_mod(vis = "pub"))]
48pub struct ConnectionOptions {
49 #[builder(start_fn, into)]
51 pub target: Url,
52 #[builder(default)]
54 pub identity: String,
55 pub metrics_meter: Option<TemporalMeter>,
58 pub tls_options: Option<TlsOptions>,
62 pub override_origin: Option<Uri>,
68 pub api_key: Option<String>,
71 pub connect_timeout: Option<Duration>,
76 #[builder(default)]
78 pub retry_options: RetryOptions,
79 #[builder(required, default = Some(ClientKeepAliveOptions::default()))]
82 pub keep_alive: Option<ClientKeepAliveOptions>,
83 pub headers: Option<HashMap<String, String>>,
89 pub binary_headers: Option<HashMap<String, Vec<u8>>>,
94 pub http_connect_proxy: Option<HttpConnectProxyOptions>,
96 #[builder(required, default = Some(DnsLoadBalancingOptions::default()))]
101 pub dns_load_balancing: Option<DnsLoadBalancingOptions>,
102 #[builder(default)]
104 pub disable_error_code_metric_tags: bool,
105 pub service_override: Option<callback_based::CallbackBasedGrpcService>,
107 #[builder(default)]
112 pub grpc_compression: GrpcCompression,
113 #[cfg(feature = "experimental")]
117 #[cfg_attr(
118 docsrs,
119 builder(setters(
120 some_fn(name = payload_limits_impl, vis = "pub(crate)"),
121 option_fn(name = maybe_payload_limits_impl, vis = "pub(crate)")
122 ))
123 )]
124 #[builder(default)]
125 pub payload_limits: PayloadLimitsOptions,
126
127 #[builder(default)]
130 #[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
131 pub(crate) skip_get_system_info: bool,
132 #[builder(default = "temporal-rust".to_owned())]
135 #[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
136 pub(crate) client_name: String,
137 #[builder(default = VERSION.to_owned())]
142 #[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
143 pub(crate) client_version: String,
144}
145
146#[cfg(all(feature = "experimental", docsrs))]
149impl<S: connection_options_builder::State> ConnectionOptionsBuilder<S> {
150 #[doc(cfg(feature = "experimental"))]
152 pub fn payload_limits(
153 self,
154 value: PayloadLimitsOptions,
155 ) -> ConnectionOptionsBuilder<connection_options_builder::SetPayloadLimits<S>>
156 where
157 S::PayloadLimits: connection_options_builder::IsUnset,
158 {
159 self.payload_limits_impl(value)
160 }
161
162 #[doc(cfg(feature = "experimental"))]
164 pub fn maybe_payload_limits(
165 self,
166 value: Option<PayloadLimitsOptions>,
167 ) -> ConnectionOptionsBuilder<connection_options_builder::SetPayloadLimits<S>>
168 where
169 S::PayloadLimits: connection_options_builder::IsUnset,
170 {
171 self.maybe_payload_limits_impl(value)
172 }
173}
174
175#[cfg(feature = "core-based-sdk")]
177impl ConnectionOptions {
178 pub fn set_skip_get_system_info(&mut self, skip: bool) {
180 self.skip_get_system_info = skip;
181 }
182 pub fn get_skip_get_system_info(&self) -> bool {
184 self.skip_get_system_info
185 }
186 pub fn get_client_name(&self) -> &str {
188 &self.client_name
189 }
190 pub fn get_client_version(&self) -> &str {
192 &self.client_version
193 }
194}
195
196#[derive(Clone, derive_more::Debug, bon::Builder)]
198#[non_exhaustive]
199#[builder(start_fn = new, on(String, into), state_mod(vis = "pub"))]
200pub struct ClientOptions {
201 #[builder(start_fn)]
203 pub namespace: String,
204
205 #[builder(field)]
206 #[debug(skip)]
207 #[cfg(feature = "experimental")]
208 plugins: Vec<ErasedClientPlugin>,
209
210 #[builder(field)]
211 #[debug(skip)]
212 #[cfg(feature = "experimental")]
213 client_plugins_applied: bool,
214
215 #[builder(default)]
217 pub data_converter: DataConverter,
218 #[builder(default)]
220 #[debug(skip)]
221 pub client_interceptors: Vec<Arc<dyn ClientInterceptor>>,
222}
223
224#[cfg(feature = "experimental")]
225impl<S: client_options_builder::State> ClientOptionsBuilder<S> {
226 pub fn plugin<P: Into<ErasedClientPlugin>>(mut self, plugin: P) -> Self {
230 self.plugins.push(plugin.into());
231 self
232 }
233
234 pub fn plugins<I, P>(mut self, plugins: I) -> Self
238 where
239 I: IntoIterator<Item = P>,
240 P: Into<ErasedClientPlugin>,
241 {
242 self.plugins.extend(plugins.into_iter().map(Into::into));
243 self
244 }
245
246 pub fn client_plugin<P: ClientPlugin>(mut self, plugin: P) -> Self {
250 self.plugins.push(ErasedClientPlugin::new(plugin));
251 self
252 }
253}
254
255impl ClientOptions {
256 #[cfg(feature = "experimental")]
262 pub fn plugins(&self) -> &[ErasedClientPlugin] {
263 &self.plugins
264 }
265
266 #[cfg(feature = "experimental")]
267 pub(crate) fn client_plugins_applied(&self) -> bool {
268 self.client_plugins_applied
269 }
270
271 #[cfg(feature = "experimental")]
272 pub(crate) fn mark_client_plugins_applied(&mut self) {
273 self.client_plugins_applied = true;
274 }
275}
276
277#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
280#[non_exhaustive]
281pub enum GrpcCompression {
282 None,
284 #[default]
286 Gzip,
287}
288
289#[derive(Clone, bon::Builder)]
291#[non_exhaustive]
292pub struct TlsOptions {
293 pub server_root_ca_cert: Option<Vec<u8>>,
297 pub domain: Option<String>,
300 pub client_tls_options: Option<ClientTlsOptions>,
305 pub server_cert_verifier: Option<Arc<dyn ServerCertVerifier>>,
322 #[cfg(feature = "dynamic-tls")]
327 pub client_cert_resolver: Option<Arc<dyn ResolvesClientCert>>,
328}
329
330impl Default for TlsOptions {
331 fn default() -> Self {
332 Self::builder().build()
333 }
334}
335
336impl std::fmt::Debug for TlsOptions {
337 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
338 let mut s = f.debug_struct("TlsOptions");
339 s.field(
340 "server_root_ca_cert",
341 &self
342 .server_root_ca_cert
343 .as_ref()
344 .map(|c| format!("{} bytes", c.len())),
345 );
346 s.field("domain", &self.domain);
347 s.field("client_tls_options", &self.client_tls_options);
348 s.field(
349 "server_cert_verifier",
350 &self.server_cert_verifier.as_ref().map(|_| "<custom>"),
351 );
352 #[cfg(feature = "dynamic-tls")]
353 s.field(
354 "client_cert_resolver",
355 &self.client_cert_resolver.as_ref().map(|_| "<custom>"),
356 );
357 s.finish()
358 }
359}
360
361#[derive(Clone, bon::Builder)]
363#[non_exhaustive]
364pub struct ClientTlsOptions {
365 pub client_cert: Vec<u8>,
367 pub client_private_key: Vec<u8>,
369}
370
371#[derive(Clone, Debug, PartialEq, bon::Builder)]
373#[non_exhaustive]
374pub struct ClientKeepAliveOptions {
375 #[builder(default = Duration::from_secs(30))]
377 pub interval: Duration,
378 #[builder(default = Duration::from_secs(15))]
380 pub timeout: Duration,
381}
382
383impl Default for ClientKeepAliveOptions {
384 fn default() -> Self {
385 Self::builder().build()
386 }
387}
388
389#[derive(Clone, Debug, PartialEq, bon::Builder)]
391#[non_exhaustive]
392pub struct DnsLoadBalancingOptions {
393 #[builder(default = Duration::from_secs(30))]
395 pub resolution_interval: Duration,
396}
397
398impl Default for DnsLoadBalancingOptions {
399 fn default() -> Self {
400 Self::builder().build()
401 }
402}
403
404#[cfg(feature = "experimental")]
407#[derive(Clone, Debug, PartialEq, bon::Builder)]
408#[non_exhaustive]
409pub struct PayloadLimitsOptions {
410 #[builder(default = DEFAULT_PAYLOADS_WARN_SIZE)]
413 pub payloads_warn_size: u64,
414 #[builder(default = DEFAULT_MEMO_WARN_SIZE)]
417 pub memo_warn_size: u64,
418}
419
420#[cfg(feature = "experimental")]
421impl Default for PayloadLimitsOptions {
422 fn default() -> Self {
423 Self::builder().build()
424 }
425}
426
427impl std::fmt::Debug for ClientTlsOptions {
428 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
430 write!(f, "ClientTlsOptions(..)")
431 }
432}
433
434#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
436#[non_exhaustive]
437pub enum WorkflowIdReusePolicy {
438 #[default]
440 Unspecified,
441 AllowDuplicate,
443 AllowDuplicateFailedOnly,
445 RejectDuplicate,
447}
448
449impl From<WorkflowIdReusePolicy> for ProtoWorkflowIdReusePolicy {
450 fn from(value: WorkflowIdReusePolicy) -> Self {
451 match value {
452 WorkflowIdReusePolicy::Unspecified => Self::Unspecified,
453 WorkflowIdReusePolicy::AllowDuplicate => Self::AllowDuplicate,
454 WorkflowIdReusePolicy::AllowDuplicateFailedOnly => Self::AllowDuplicateFailedOnly,
455 WorkflowIdReusePolicy::RejectDuplicate => Self::RejectDuplicate,
456 }
457 }
458}
459
460#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
463#[non_exhaustive]
464pub enum WorkflowIdConflictPolicy {
465 #[default]
467 Unspecified,
468 Fail,
470 UseExisting,
472 TerminateExisting,
474}
475
476impl From<WorkflowIdConflictPolicy> for ProtoWorkflowIdConflictPolicy {
477 fn from(value: WorkflowIdConflictPolicy) -> Self {
478 match value {
479 WorkflowIdConflictPolicy::Unspecified => Self::Unspecified,
480 WorkflowIdConflictPolicy::Fail => Self::Fail,
481 WorkflowIdConflictPolicy::UseExisting => Self::UseExisting,
482 WorkflowIdConflictPolicy::TerminateExisting => Self::TerminateExisting,
483 }
484 }
485}
486
487#[derive(Debug, Clone, bon::Builder)]
489#[builder(start_fn = new, on(String, into))]
490#[non_exhaustive]
491pub struct WorkflowStartOptions {
492 #[builder(start_fn)]
494 pub task_queue: String,
495
496 #[builder(start_fn)]
498 pub workflow_id: String,
499
500 #[builder(default)]
502 pub id_reuse_policy: WorkflowIdReusePolicy,
503
504 #[builder(default)]
507 pub id_conflict_policy: WorkflowIdConflictPolicy,
508
509 pub execution_timeout: Option<Duration>,
512
513 pub run_timeout: Option<Duration>,
515
516 pub task_timeout: Option<Duration>,
518
519 pub cron_schedule: Option<String>,
521
522 pub search_attributes: Option<SearchAttributes>,
524
525 #[builder(default)]
527 pub enable_eager_workflow_start: bool,
528
529 #[builder(into)]
531 pub retry_policy: Option<RetryPolicy>,
532
533 #[builder(default)]
535 pub links: Vec<common::v1::Link>,
536
537 #[builder(default)]
540 pub completion_callbacks: Vec<common::v1::Callback>,
541
542 #[builder(default)]
544 pub priority: Priority,
545
546 pub header: Option<Header>,
548
549 pub memo: Option<MemoValues>,
551
552 pub static_summary: Option<String>,
554
555 pub static_details: Option<String>,
557
558 #[builder(default)]
560 pub rpc_options: RpcOptions,
561}
562
563impl WorkflowStartOptions {
564 pub(crate) async fn encoded_memo(
565 &self,
566 data_converter: &DataConverter,
567 ) -> Result<Option<ProtoMemo>, PayloadConversionError> {
568 let Some(memo) = &self.memo else {
569 return Ok(None);
570 };
571
572 let payload_converter = data_converter.payload_converter();
573 let context_data = SerializationContextData::Workflow(WorkflowSerializationContext::new());
574 let context = SerializationContext::new(&context_data, payload_converter);
575 let mut memo = ProtoMemo {
576 fields: memo
577 .iter()
578 .map(|(key, value)| {
579 payload_converter
580 .to_payload(&context, value)
581 .map(|payload| (key.to_owned(), payload))
582 })
583 .collect::<Result<_, _>>()?,
584 };
585 encode_payloads(
586 &mut memo,
587 data_converter.codec(),
588 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
589 )
590 .await?;
591 Ok(Some(memo))
592 }
593
594 pub(crate) fn user_metadata(&self) -> Option<UserMetadata> {
595 (self.static_summary.is_some() || self.static_details.is_some()).then(|| {
596 let payload_converter = PayloadConverter::default();
597 let context_data =
598 SerializationContextData::Workflow(WorkflowSerializationContext::new());
599 let context = SerializationContext::new(&context_data, &payload_converter);
600 UserMetadata {
601 summary: self.static_summary.as_ref().map(|summary| {
602 payload_converter
603 .to_payload(&context, summary)
604 .expect("String-to-JSON payload serialization is infallible")
605 }),
606 details: self.static_details.as_ref().map(|details| {
607 payload_converter
608 .to_payload(&context, details)
609 .expect("String-to-JSON payload serialization is infallible")
610 }),
611 }
612 })
613 }
614}
615
616#[derive(Debug, Clone, bon::Builder)]
621#[builder(start_fn = new, on(String, into))]
622#[non_exhaustive]
623pub struct WorkflowUpdateWithStartOptions {
624 #[builder(start_fn)]
626 pub task_queue: String,
627
628 #[builder(start_fn)]
630 pub workflow_id: String,
631
632 #[builder(start_fn)]
635 pub id_conflict_policy: WorkflowIdConflictPolicy,
636
637 #[builder(default)]
639 pub id_reuse_policy: WorkflowIdReusePolicy,
640
641 pub execution_timeout: Option<Duration>,
643
644 pub run_timeout: Option<Duration>,
646
647 pub task_timeout: Option<Duration>,
649
650 pub search_attributes: Option<SearchAttributes>,
652
653 #[builder(into)]
655 pub retry_policy: Option<RetryPolicy>,
656
657 #[builder(default)]
659 pub links: Vec<common::v1::Link>,
660
661 #[builder(default)]
663 pub completion_callbacks: Vec<common::v1::Callback>,
664
665 #[builder(default)]
667 pub priority: Priority,
668
669 pub start_header: Option<Header>,
671
672 pub update_header: Option<Header>,
674
675 pub memo: Option<MemoValues>,
677
678 pub static_summary: Option<String>,
680
681 pub static_details: Option<String>,
683
684 pub update_id: Option<String>,
686
687 #[builder(default)]
689 pub rpc_options: RpcOptions,
690}
691
692impl WorkflowUpdateWithStartOptions {
693 pub(crate) fn into_parts(self) -> (WorkflowStartOptions, Option<String>, Option<Header>) {
694 let Self {
695 task_queue,
696 workflow_id,
697 id_conflict_policy,
698 id_reuse_policy,
699 execution_timeout,
700 run_timeout,
701 task_timeout,
702 search_attributes,
703 retry_policy,
704 links,
705 completion_callbacks,
706 priority,
707 start_header,
708 update_header,
709 memo,
710 static_summary,
711 static_details,
712 update_id,
713 rpc_options: _,
714 } = self;
715 (
716 WorkflowStartOptions {
717 task_queue,
718 workflow_id,
719 id_reuse_policy,
720 id_conflict_policy,
721 execution_timeout,
722 run_timeout,
723 task_timeout,
724 cron_schedule: None,
725 search_attributes,
726 enable_eager_workflow_start: false,
727 retry_policy,
728 links,
729 completion_callbacks,
730 priority,
731 header: start_header,
732 memo,
733 static_summary,
734 static_details,
735 rpc_options: RpcOptions::default(),
736 },
737 update_id,
738 update_header,
739 )
740 }
741}
742
743pub use temporalio_common::Priority;
744
745#[derive(Debug, Clone, bon::Builder)]
747#[non_exhaustive]
748pub struct WorkflowGetResultOptions {
749 #[builder(default = true)]
752 pub follow_runs: bool,
753 #[builder(default)]
755 pub rpc_options: RpcOptions,
756}
757impl Default for WorkflowGetResultOptions {
758 fn default() -> Self {
759 Self {
760 follow_runs: true,
761 rpc_options: RpcOptions::default(),
762 }
763 }
764}
765
766#[derive(Debug, Clone, Default, bon::Builder)]
768#[non_exhaustive]
769pub struct WorkflowExecuteUpdateOptions {
770 pub update_id: Option<String>,
772 pub header: Option<Header>,
774 #[builder(default)]
776 pub rpc_options: RpcOptions,
777}
778
779#[derive(Debug, Clone, Default, bon::Builder)]
781#[non_exhaustive]
782pub struct WorkflowSignalOptions {
783 pub request_id: Option<String>,
785 pub header: Option<Header>,
787 #[builder(default)]
789 pub rpc_options: RpcOptions,
790}
791
792#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
794#[non_exhaustive]
795pub enum QueryRejectCondition {
796 #[default]
798 Unspecified,
799 None,
801 NotOpen,
803 NotCompletedCleanly,
805}
806
807impl From<QueryRejectCondition> for ProtoQueryRejectCondition {
808 fn from(value: QueryRejectCondition) -> Self {
809 match value {
810 QueryRejectCondition::Unspecified => Self::Unspecified,
811 QueryRejectCondition::None => Self::None,
812 QueryRejectCondition::NotOpen => Self::NotOpen,
813 QueryRejectCondition::NotCompletedCleanly => Self::NotCompletedCleanly,
814 }
815 }
816}
817
818#[derive(Debug, Clone, Default, bon::Builder)]
820#[non_exhaustive]
821pub struct WorkflowQueryOptions {
822 pub reject_condition: Option<QueryRejectCondition>,
825 pub header: Option<Header>,
827 #[builder(default)]
829 pub rpc_options: RpcOptions,
830}
831
832#[derive(Debug, Clone, Default, bon::Builder)]
834#[builder(on(String, into))]
835#[non_exhaustive]
836pub struct WorkflowCancelOptions {
837 #[builder(default)]
839 pub reason: String,
840 pub request_id: Option<String>,
842 #[builder(default)]
844 pub rpc_options: RpcOptions,
845}
846
847#[derive(Debug, Clone, Default, bon::Builder)]
849#[builder(on(String, into))]
850#[non_exhaustive]
851pub struct WorkflowTerminateOptions {
852 #[builder(default)]
854 pub reason: String,
855 pub details: Option<Payloads>,
857 #[builder(default)]
859 pub rpc_options: RpcOptions,
860}
861
862#[derive(Debug, Clone, Default, bon::Builder)]
864#[non_exhaustive]
865pub struct WorkflowDescribeOptions {
866 #[builder(default)]
868 pub rpc_options: RpcOptions,
869}
870
871const DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD: Duration = Duration::from_secs(60 * 60 * 24 * 3);
873
874#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
876#[non_exhaustive]
877pub enum ArchivalState {
878 #[default]
880 Unspecified,
881 Disabled,
883 Enabled,
885}
886
887impl From<ArchivalState> for ProtoArchivalState {
888 fn from(value: ArchivalState) -> Self {
889 match value {
890 ArchivalState::Unspecified => Self::Unspecified,
891 ArchivalState::Disabled => Self::Disabled,
892 ArchivalState::Enabled => Self::Enabled,
893 }
894 }
895}
896
897#[derive(Clone, Debug, bon::Builder)]
899#[builder(on(String, into))]
900#[non_exhaustive]
901pub struct RegisterNamespaceOptions {
902 pub namespace: String,
904 pub description: String,
906 #[builder(default)]
908 pub owner_email: String,
909 #[builder(default = DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD)]
911 pub workflow_execution_retention_period: Duration,
912 #[builder(default)]
914 pub clusters: Vec<ClusterReplicationConfig>,
915 #[builder(default)]
917 pub active_cluster_name: String,
918 #[builder(default)]
920 pub data: HashMap<String, String>,
921 #[builder(default)]
923 pub security_token: String,
924 #[builder(default)]
926 pub is_global_namespace: bool,
927 #[builder(default = ArchivalState::Unspecified)]
929 pub history_archival_state: ArchivalState,
930 #[builder(default)]
932 pub history_archival_uri: String,
933 #[builder(default = ArchivalState::Unspecified)]
935 pub visibility_archival_state: ArchivalState,
936 #[builder(default)]
938 pub visibility_archival_uri: String,
939}
940
941impl From<RegisterNamespaceOptions> for RegisterNamespaceRequest {
942 fn from(val: RegisterNamespaceOptions) -> Self {
943 RegisterNamespaceRequest {
944 namespace: val.namespace,
945 description: val.description,
946 owner_email: val.owner_email,
947 workflow_execution_retention_period: val
948 .workflow_execution_retention_period
949 .try_into()
950 .ok(),
951 clusters: val.clusters,
952 active_cluster_name: val.active_cluster_name,
953 data: val.data,
954 security_token: val.security_token,
955 is_global_namespace: val.is_global_namespace,
956 history_archival_state: ProtoArchivalState::from(val.history_archival_state) as i32,
957 history_archival_uri: val.history_archival_uri,
958 visibility_archival_state: ProtoArchivalState::from(val.visibility_archival_state)
959 as i32,
960 visibility_archival_uri: val.visibility_archival_uri,
961 }
962 }
963}
964
965#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
967#[non_exhaustive]
968pub enum HistoryEventFilterType {
969 #[default]
971 Unspecified,
972 AllEvent,
974 CloseEvent,
976}
977
978impl From<HistoryEventFilterType> for ProtoHistoryEventFilterType {
979 fn from(value: HistoryEventFilterType) -> Self {
980 match value {
981 HistoryEventFilterType::Unspecified => Self::Unspecified,
982 HistoryEventFilterType::AllEvent => Self::AllEvent,
983 HistoryEventFilterType::CloseEvent => Self::CloseEvent,
984 }
985 }
986}
987
988#[derive(Debug, Clone, Default, bon::Builder)]
990#[non_exhaustive]
991pub struct WorkflowFetchHistoryOptions {
992 #[builder(default)]
994 pub skip_archival: bool,
995 #[builder(default)]
997 pub wait_new_event: bool,
998 #[builder(default = HistoryEventFilterType::AllEvent)]
1000 pub event_filter_type: HistoryEventFilterType,
1001 #[builder(default)]
1003 pub rpc_options: RpcOptions,
1004}
1005
1006#[derive(Debug, Clone, Default, bon::Builder)]
1008#[non_exhaustive]
1009pub struct WorkflowStartUpdateOptions {
1010 pub update_id: Option<String>,
1012 pub header: Option<Header>,
1014 #[builder(default)]
1016 pub rpc_options: RpcOptions,
1017}
1018
1019impl From<WorkflowExecuteUpdateOptions> for WorkflowStartUpdateOptions {
1020 fn from(options: WorkflowExecuteUpdateOptions) -> Self {
1022 Self::builder()
1023 .maybe_update_id(options.update_id)
1024 .maybe_header(options.header)
1025 .rpc_options(options.rpc_options)
1026 .build()
1027 }
1028}
1029
1030#[derive(Debug, Clone, Default, bon::Builder)]
1032#[non_exhaustive]
1033pub struct WorkflowListOptions {
1034 pub limit: Option<usize>,
1037 #[builder(default)]
1039 pub rpc_options: RpcOptions,
1040}
1041
1042#[derive(Debug, Clone, Default, bon::Builder)]
1044#[non_exhaustive]
1045pub struct WorkflowCountOptions {
1046 #[builder(default)]
1048 pub rpc_options: RpcOptions,
1049}
1050
1051#[derive(Clone, Debug, bon::Builder)]
1053#[builder(start_fn = new, on(String, into))]
1054#[non_exhaustive]
1055pub struct ActivityStartOptions {
1056 #[builder(start_fn)]
1058 pub task_queue: String,
1059 #[builder(start_fn)]
1061 pub id: String,
1062 #[builder(start_fn)]
1066 pub close_timeouts: ActivityCloseTimeouts,
1067 pub schedule_to_start_timeout: Option<Duration>,
1070 pub heartbeat_timeout: Option<Duration>,
1072 #[builder(into)]
1074 pub retry_policy: Option<RetryPolicy>,
1075 #[builder(default)]
1077 pub priority: Priority,
1078 #[builder(default)]
1080 pub id_reuse_policy: ActivityIdReusePolicy,
1081 #[builder(default)]
1084 pub id_conflict_policy: ActivityIdConflictPolicy,
1085 pub search_attributes: Option<SearchAttributes>,
1087 pub header: Option<Header>,
1089 pub summary: Option<String>,
1091 pub static_details: Option<String>,
1093 pub start_delay: Option<Duration>,
1096}
1097
1098impl ActivityStartOptions {
1099 pub fn with_start_to_close_timeout(
1101 task_queue: impl Into<String>,
1102 activity_id: impl Into<String>,
1103 start_to_close_timeout: Duration,
1104 ) -> ActivityStartOptionsBuilder {
1105 Self::new(
1106 task_queue,
1107 activity_id,
1108 ActivityCloseTimeouts::StartToClose(start_to_close_timeout),
1109 )
1110 }
1111
1112 pub fn with_schedule_to_close_timeout(
1114 task_queue: impl Into<String>,
1115 activity_id: impl Into<String>,
1116 schedule_to_close_timeout: Duration,
1117 ) -> ActivityStartOptionsBuilder {
1118 Self::new(
1119 task_queue,
1120 activity_id,
1121 ActivityCloseTimeouts::ScheduleToClose(schedule_to_close_timeout),
1122 )
1123 }
1124}
1125
1126#[non_exhaustive]
1129#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
1130pub enum ActivityIdReusePolicy {
1131 #[default]
1132 AllowDuplicate,
1134 AllowDuplicateFailedOnly,
1137 RejectDuplicate,
1139}
1140
1141impl From<ActivityIdReusePolicy> for ProtoActivityIdReusePolicy {
1142 fn from(value: ActivityIdReusePolicy) -> Self {
1143 match value {
1144 ActivityIdReusePolicy::AllowDuplicate => Self::AllowDuplicate,
1145 ActivityIdReusePolicy::AllowDuplicateFailedOnly => Self::AllowDuplicateFailedOnly,
1146 ActivityIdReusePolicy::RejectDuplicate => Self::RejectDuplicate,
1147 }
1148 }
1149}
1150
1151#[non_exhaustive]
1154#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
1155pub enum ActivityIdConflictPolicy {
1156 #[default]
1157 Fail,
1160 UseExisting,
1162}
1163
1164impl From<ActivityIdConflictPolicy> for ProtoActivityIdConflictPolicy {
1165 fn from(value: ActivityIdConflictPolicy) -> Self {
1166 match value {
1167 ActivityIdConflictPolicy::Fail => Self::Fail,
1168 ActivityIdConflictPolicy::UseExisting => Self::UseExisting,
1169 }
1170 }
1171}
1172
1173#[derive(Debug, Clone, Default, bon::Builder)]
1175#[non_exhaustive]
1176pub struct ActivityListOptions {}
1177
1178#[derive(Debug, Clone, Default, bon::Builder)]
1180#[non_exhaustive]
1181pub struct ActivityCountOptions {}
1182
1183#[derive(Debug, Clone, Default, bon::Builder)]
1191#[non_exhaustive]
1192pub struct ActivityDescribeOptions {
1193 #[builder(default)]
1195 pub include_input: bool,
1196 #[builder(default)]
1198 pub include_outcome: bool,
1199 #[builder(default)]
1201 pub include_heartbeat_details: bool,
1202 #[builder(default)]
1204 pub include_last_failure: bool,
1205}
1206
1207#[derive(Debug, Clone, Default, bon::Builder)]
1209#[builder(on(String, into))]
1210#[non_exhaustive]
1211pub struct ActivityCancelOptions {
1212 #[builder(default)]
1214 pub reason: String,
1215}
1216
1217#[derive(Debug, Clone, Default, bon::Builder)]
1219#[builder(on(String, into))]
1220#[non_exhaustive]
1221pub struct ActivityTerminateOptions {
1222 #[builder(default)]
1224 pub reason: String,
1225}