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, VersioningOverride,
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 versioning_override: Option<VersioningOverride>,
551
552 pub header: Option<Header>,
554
555 pub memo: Option<MemoValues>,
557
558 pub static_summary: Option<String>,
560
561 pub static_details: Option<String>,
563
564 #[builder(default)]
566 pub rpc_options: RpcOptions,
567}
568
569impl WorkflowStartOptions {
570 pub(crate) async fn encoded_memo(
571 &self,
572 data_converter: &DataConverter,
573 ) -> Result<Option<ProtoMemo>, PayloadConversionError> {
574 let Some(memo) = &self.memo else {
575 return Ok(None);
576 };
577
578 let payload_converter = data_converter.payload_converter();
579 let context_data = SerializationContextData::Workflow(WorkflowSerializationContext::new());
580 let context = SerializationContext::new(&context_data, payload_converter);
581 let mut memo = ProtoMemo {
582 fields: memo
583 .iter()
584 .map(|(key, value)| {
585 payload_converter
586 .to_payload(&context, value)
587 .map(|payload| (key.to_owned(), payload))
588 })
589 .collect::<Result<_, _>>()?,
590 };
591 encode_payloads(
592 &mut memo,
593 data_converter.codec(),
594 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
595 )
596 .await?;
597 Ok(Some(memo))
598 }
599
600 pub(crate) fn user_metadata(&self) -> Option<UserMetadata> {
601 (self.static_summary.is_some() || self.static_details.is_some()).then(|| {
602 let payload_converter = PayloadConverter::default();
603 let context_data =
604 SerializationContextData::Workflow(WorkflowSerializationContext::new());
605 let context = SerializationContext::new(&context_data, &payload_converter);
606 UserMetadata {
607 summary: self.static_summary.as_ref().map(|summary| {
608 payload_converter
609 .to_payload(&context, summary)
610 .expect("String-to-JSON payload serialization is infallible")
611 }),
612 details: self.static_details.as_ref().map(|details| {
613 payload_converter
614 .to_payload(&context, details)
615 .expect("String-to-JSON payload serialization is infallible")
616 }),
617 }
618 })
619 }
620}
621
622#[derive(Debug, Clone, bon::Builder)]
627#[builder(start_fn = new, on(String, into))]
628#[non_exhaustive]
629pub struct WorkflowUpdateWithStartOptions {
630 #[builder(start_fn)]
632 pub task_queue: String,
633
634 #[builder(start_fn)]
636 pub workflow_id: String,
637
638 #[builder(start_fn)]
641 pub id_conflict_policy: WorkflowIdConflictPolicy,
642
643 #[builder(default)]
645 pub id_reuse_policy: WorkflowIdReusePolicy,
646
647 pub execution_timeout: Option<Duration>,
649
650 pub run_timeout: Option<Duration>,
652
653 pub task_timeout: Option<Duration>,
655
656 pub search_attributes: Option<SearchAttributes>,
658
659 #[builder(into)]
661 pub retry_policy: Option<RetryPolicy>,
662
663 #[builder(default)]
665 pub links: Vec<common::v1::Link>,
666
667 #[builder(default)]
669 pub completion_callbacks: Vec<common::v1::Callback>,
670
671 #[builder(default)]
673 pub priority: Priority,
674
675 pub versioning_override: Option<VersioningOverride>,
680
681 pub start_header: Option<Header>,
683
684 pub update_header: Option<Header>,
686
687 pub memo: Option<MemoValues>,
689
690 pub static_summary: Option<String>,
692
693 pub static_details: Option<String>,
695
696 pub update_id: Option<String>,
698
699 #[builder(default)]
701 pub rpc_options: RpcOptions,
702}
703
704impl WorkflowUpdateWithStartOptions {
705 pub(crate) fn into_parts(self) -> (WorkflowStartOptions, Option<String>, Option<Header>) {
706 let Self {
707 task_queue,
708 workflow_id,
709 id_conflict_policy,
710 id_reuse_policy,
711 execution_timeout,
712 run_timeout,
713 task_timeout,
714 search_attributes,
715 retry_policy,
716 links,
717 completion_callbacks,
718 priority,
719 versioning_override,
720 start_header,
721 update_header,
722 memo,
723 static_summary,
724 static_details,
725 update_id,
726 rpc_options: _,
727 } = self;
728 (
729 WorkflowStartOptions {
730 task_queue,
731 workflow_id,
732 id_reuse_policy,
733 id_conflict_policy,
734 execution_timeout,
735 run_timeout,
736 task_timeout,
737 cron_schedule: None,
738 search_attributes,
739 enable_eager_workflow_start: false,
740 retry_policy,
741 links,
742 completion_callbacks,
743 priority,
744 versioning_override,
745 header: start_header,
746 memo,
747 static_summary,
748 static_details,
749 rpc_options: RpcOptions::default(),
750 },
751 update_id,
752 update_header,
753 )
754 }
755}
756
757pub use temporalio_common::Priority;
758
759#[derive(Debug, Clone, bon::Builder)]
761#[non_exhaustive]
762pub struct WorkflowGetResultOptions {
763 #[builder(default = true)]
766 pub follow_runs: bool,
767 #[builder(default)]
769 pub rpc_options: RpcOptions,
770}
771impl Default for WorkflowGetResultOptions {
772 fn default() -> Self {
773 Self {
774 follow_runs: true,
775 rpc_options: RpcOptions::default(),
776 }
777 }
778}
779
780#[derive(Debug, Clone, Default, bon::Builder)]
782#[non_exhaustive]
783pub struct WorkflowExecuteUpdateOptions {
784 pub update_id: Option<String>,
786 pub header: Option<Header>,
788 #[builder(default)]
790 pub rpc_options: RpcOptions,
791}
792
793#[derive(Debug, Clone, Default, bon::Builder)]
795#[non_exhaustive]
796pub struct WorkflowSignalOptions {
797 pub request_id: Option<String>,
799 pub header: Option<Header>,
801 #[builder(default)]
803 pub rpc_options: RpcOptions,
804}
805
806#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
808#[non_exhaustive]
809pub enum QueryRejectCondition {
810 #[default]
812 Unspecified,
813 None,
815 NotOpen,
817 NotCompletedCleanly,
819}
820
821impl From<QueryRejectCondition> for ProtoQueryRejectCondition {
822 fn from(value: QueryRejectCondition) -> Self {
823 match value {
824 QueryRejectCondition::Unspecified => Self::Unspecified,
825 QueryRejectCondition::None => Self::None,
826 QueryRejectCondition::NotOpen => Self::NotOpen,
827 QueryRejectCondition::NotCompletedCleanly => Self::NotCompletedCleanly,
828 }
829 }
830}
831
832#[derive(Debug, Clone, Default, bon::Builder)]
834#[non_exhaustive]
835pub struct WorkflowQueryOptions {
836 pub reject_condition: Option<QueryRejectCondition>,
839 pub header: Option<Header>,
841 #[builder(default)]
843 pub rpc_options: RpcOptions,
844}
845
846#[derive(Debug, Clone, Default, bon::Builder)]
848#[builder(on(String, into))]
849#[non_exhaustive]
850pub struct WorkflowCancelOptions {
851 #[builder(default)]
853 pub reason: String,
854 pub request_id: Option<String>,
856 #[builder(default)]
858 pub rpc_options: RpcOptions,
859}
860
861#[derive(Debug, Clone, Default, bon::Builder)]
863#[builder(on(String, into))]
864#[non_exhaustive]
865pub struct WorkflowTerminateOptions {
866 #[builder(default)]
868 pub reason: String,
869 pub details: Option<Payloads>,
871 #[builder(default)]
873 pub rpc_options: RpcOptions,
874}
875
876#[derive(Debug, Clone, Default, bon::Builder)]
878#[non_exhaustive]
879pub struct WorkflowDescribeOptions {
880 #[builder(default)]
882 pub rpc_options: RpcOptions,
883}
884
885const DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD: Duration = Duration::from_secs(60 * 60 * 24 * 3);
887
888#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
890#[non_exhaustive]
891pub enum ArchivalState {
892 #[default]
894 Unspecified,
895 Disabled,
897 Enabled,
899}
900
901impl From<ArchivalState> for ProtoArchivalState {
902 fn from(value: ArchivalState) -> Self {
903 match value {
904 ArchivalState::Unspecified => Self::Unspecified,
905 ArchivalState::Disabled => Self::Disabled,
906 ArchivalState::Enabled => Self::Enabled,
907 }
908 }
909}
910
911#[derive(Clone, Debug, bon::Builder)]
913#[builder(on(String, into))]
914#[non_exhaustive]
915pub struct RegisterNamespaceOptions {
916 pub namespace: String,
918 pub description: String,
920 #[builder(default)]
922 pub owner_email: String,
923 #[builder(default = DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD)]
925 pub workflow_execution_retention_period: Duration,
926 #[builder(default)]
928 pub clusters: Vec<ClusterReplicationConfig>,
929 #[builder(default)]
931 pub active_cluster_name: String,
932 #[builder(default)]
934 pub data: HashMap<String, String>,
935 #[builder(default)]
937 pub security_token: String,
938 #[builder(default)]
940 pub is_global_namespace: bool,
941 #[builder(default = ArchivalState::Unspecified)]
943 pub history_archival_state: ArchivalState,
944 #[builder(default)]
946 pub history_archival_uri: String,
947 #[builder(default = ArchivalState::Unspecified)]
949 pub visibility_archival_state: ArchivalState,
950 #[builder(default)]
952 pub visibility_archival_uri: String,
953}
954
955impl From<RegisterNamespaceOptions> for RegisterNamespaceRequest {
956 fn from(val: RegisterNamespaceOptions) -> Self {
957 RegisterNamespaceRequest {
958 namespace: val.namespace,
959 description: val.description,
960 owner_email: val.owner_email,
961 workflow_execution_retention_period: val
962 .workflow_execution_retention_period
963 .try_into()
964 .ok(),
965 clusters: val.clusters,
966 active_cluster_name: val.active_cluster_name,
967 data: val.data,
968 security_token: val.security_token,
969 is_global_namespace: val.is_global_namespace,
970 history_archival_state: ProtoArchivalState::from(val.history_archival_state) as i32,
971 history_archival_uri: val.history_archival_uri,
972 visibility_archival_state: ProtoArchivalState::from(val.visibility_archival_state)
973 as i32,
974 visibility_archival_uri: val.visibility_archival_uri,
975 }
976 }
977}
978
979#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
981#[non_exhaustive]
982pub enum HistoryEventFilterType {
983 #[default]
985 Unspecified,
986 AllEvent,
988 CloseEvent,
990}
991
992impl From<HistoryEventFilterType> for ProtoHistoryEventFilterType {
993 fn from(value: HistoryEventFilterType) -> Self {
994 match value {
995 HistoryEventFilterType::Unspecified => Self::Unspecified,
996 HistoryEventFilterType::AllEvent => Self::AllEvent,
997 HistoryEventFilterType::CloseEvent => Self::CloseEvent,
998 }
999 }
1000}
1001
1002#[derive(Debug, Clone, Default, bon::Builder)]
1004#[non_exhaustive]
1005pub struct WorkflowFetchHistoryOptions {
1006 #[builder(default)]
1008 pub skip_archival: bool,
1009 #[builder(default)]
1011 pub wait_new_event: bool,
1012 #[builder(default = HistoryEventFilterType::AllEvent)]
1014 pub event_filter_type: HistoryEventFilterType,
1015 #[builder(default)]
1017 pub rpc_options: RpcOptions,
1018}
1019
1020#[derive(Debug, Clone, Default, bon::Builder)]
1022#[non_exhaustive]
1023pub struct WorkflowStartUpdateOptions {
1024 pub update_id: Option<String>,
1026 pub header: Option<Header>,
1028 #[builder(default)]
1030 pub rpc_options: RpcOptions,
1031}
1032
1033impl From<WorkflowExecuteUpdateOptions> for WorkflowStartUpdateOptions {
1034 fn from(options: WorkflowExecuteUpdateOptions) -> Self {
1036 Self::builder()
1037 .maybe_update_id(options.update_id)
1038 .maybe_header(options.header)
1039 .rpc_options(options.rpc_options)
1040 .build()
1041 }
1042}
1043
1044#[derive(Debug, Clone, Default, bon::Builder)]
1046#[non_exhaustive]
1047pub struct WorkflowListOptions {
1048 pub limit: Option<usize>,
1051 #[builder(default)]
1053 pub rpc_options: RpcOptions,
1054}
1055
1056#[derive(Debug, Clone, Default, bon::Builder)]
1058#[non_exhaustive]
1059pub struct WorkflowCountOptions {
1060 #[builder(default)]
1062 pub rpc_options: RpcOptions,
1063}
1064
1065#[derive(Clone, Debug, bon::Builder)]
1067#[builder(start_fn = new, on(String, into))]
1068#[non_exhaustive]
1069pub struct ActivityStartOptions {
1070 #[builder(start_fn)]
1072 pub task_queue: String,
1073 #[builder(start_fn)]
1075 pub id: String,
1076 #[builder(start_fn)]
1080 pub close_timeouts: ActivityCloseTimeouts,
1081 pub schedule_to_start_timeout: Option<Duration>,
1084 pub heartbeat_timeout: Option<Duration>,
1086 #[builder(into)]
1088 pub retry_policy: Option<RetryPolicy>,
1089 #[builder(default)]
1091 pub priority: Priority,
1092 #[builder(default)]
1094 pub id_reuse_policy: ActivityIdReusePolicy,
1095 #[builder(default)]
1098 pub id_conflict_policy: ActivityIdConflictPolicy,
1099 pub search_attributes: Option<SearchAttributes>,
1101 pub header: Option<Header>,
1103 pub summary: Option<String>,
1105 pub static_details: Option<String>,
1107 pub start_delay: Option<Duration>,
1110}
1111
1112impl ActivityStartOptions {
1113 pub fn with_start_to_close_timeout(
1115 task_queue: impl Into<String>,
1116 activity_id: impl Into<String>,
1117 start_to_close_timeout: Duration,
1118 ) -> ActivityStartOptionsBuilder {
1119 Self::new(
1120 task_queue,
1121 activity_id,
1122 ActivityCloseTimeouts::StartToClose(start_to_close_timeout),
1123 )
1124 }
1125
1126 pub fn with_schedule_to_close_timeout(
1128 task_queue: impl Into<String>,
1129 activity_id: impl Into<String>,
1130 schedule_to_close_timeout: Duration,
1131 ) -> ActivityStartOptionsBuilder {
1132 Self::new(
1133 task_queue,
1134 activity_id,
1135 ActivityCloseTimeouts::ScheduleToClose(schedule_to_close_timeout),
1136 )
1137 }
1138}
1139
1140#[non_exhaustive]
1143#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
1144pub enum ActivityIdReusePolicy {
1145 #[default]
1146 AllowDuplicate,
1148 AllowDuplicateFailedOnly,
1151 RejectDuplicate,
1153}
1154
1155impl From<ActivityIdReusePolicy> for ProtoActivityIdReusePolicy {
1156 fn from(value: ActivityIdReusePolicy) -> Self {
1157 match value {
1158 ActivityIdReusePolicy::AllowDuplicate => Self::AllowDuplicate,
1159 ActivityIdReusePolicy::AllowDuplicateFailedOnly => Self::AllowDuplicateFailedOnly,
1160 ActivityIdReusePolicy::RejectDuplicate => Self::RejectDuplicate,
1161 }
1162 }
1163}
1164
1165#[non_exhaustive]
1168#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
1169pub enum ActivityIdConflictPolicy {
1170 #[default]
1171 Fail,
1174 UseExisting,
1176}
1177
1178impl From<ActivityIdConflictPolicy> for ProtoActivityIdConflictPolicy {
1179 fn from(value: ActivityIdConflictPolicy) -> Self {
1180 match value {
1181 ActivityIdConflictPolicy::Fail => Self::Fail,
1182 ActivityIdConflictPolicy::UseExisting => Self::UseExisting,
1183 }
1184 }
1185}
1186
1187#[derive(Debug, Clone, Default, bon::Builder)]
1189#[non_exhaustive]
1190pub struct ActivityListOptions {}
1191
1192#[derive(Debug, Clone, Default, bon::Builder)]
1194#[non_exhaustive]
1195pub struct ActivityCountOptions {}
1196
1197#[derive(Debug, Clone, Default, bon::Builder)]
1205#[non_exhaustive]
1206pub struct ActivityDescribeOptions {
1207 #[builder(default)]
1209 pub include_input: bool,
1210 #[builder(default)]
1212 pub include_outcome: bool,
1213 #[builder(default)]
1215 pub include_heartbeat_details: bool,
1216 #[builder(default)]
1218 pub include_last_failure: bool,
1219}
1220
1221#[derive(Debug, Clone, Default, bon::Builder)]
1223#[builder(on(String, into))]
1224#[non_exhaustive]
1225pub struct ActivityCancelOptions {
1226 #[builder(default)]
1228 pub reason: String,
1229}
1230
1231#[derive(Debug, Clone, Default, bon::Builder)]
1233#[builder(on(String, into))]
1234#[non_exhaustive]
1235pub struct ActivityTerminateOptions {
1236 #[builder(default)]
1238 pub reason: String,
1239}