Skip to main content

temporalio_client/
options_structs.rs

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/// Options for [crate::Connection::connect].
45#[derive(bon::Builder, Clone, Debug)]
46#[non_exhaustive]
47#[builder(start_fn = new, on(String, into), state_mod(vis = "pub"))]
48pub struct ConnectionOptions {
49    /// The server to connect to.
50    #[builder(start_fn, into)]
51    pub target: Url,
52    /// A human-readable string that can identify this process. Defaults to empty string.
53    #[builder(default)]
54    pub identity: String,
55    /// When set, this client will record metrics using the provided meter. The meter can be
56    /// obtained from [temporalio_common::telemetry::TelemetryInstance::get_temporal_metric_meter].
57    pub metrics_meter: Option<TemporalMeter>,
58    /// If specified, use TLS as configured by the [TlsOptions] struct. If this is set core will
59    /// attempt to use TLS when connecting to the Temporal server. Lang SDK is expected to pass any
60    /// certs or keys as bytes, loading them from disk itself if needed.
61    pub tls_options: Option<TlsOptions>,
62    /// If set, override the origin used when connecting. May be useful in rare situations where tls
63    /// verification needs to use a different name from what should be set as the `:authority`
64    /// header. If [TlsOptions::domain] is set, and this is not, this will be set to
65    /// `https://<domain>`, effectively making the `:authority` header consistent with the domain
66    /// override.
67    pub override_origin: Option<Uri>,
68    /// An API key to use for auth. If set, TLS will be enabled by default, but without any mTLS
69    /// specific settings.
70    pub api_key: Option<String>,
71    /// When set, limits the time allowed to establish the initial TCP/TLS connection to the
72    /// server. If the connection cannot be established within this duration, `connect` will
73    /// return an error. When `None` (the default), no explicit timeout is applied and the
74    /// connection attempt may block indefinitely (subject to OS-level TCP timeouts).
75    pub connect_timeout: Option<Duration>,
76    /// Retry configuration for the server client. Default is [RetryOptions::default]
77    #[builder(default)]
78    pub retry_options: RetryOptions,
79    /// If set, HTTP2 gRPC keep alive will be enabled.
80    /// To enable with default settings, use `.keep_alive(Some(ClientKeepAliveConfig::default()))`.
81    #[builder(required, default = Some(ClientKeepAliveOptions::default()))]
82    pub keep_alive: Option<ClientKeepAliveOptions>,
83    /// HTTP headers to include on every RPC call.
84    ///
85    /// These must be valid gRPC metadata keys, and must not be binary metadata keys (ending in
86    /// `-bin). To set binary headers, use [ConnectionOptions::binary_headers]. Invalid header keys
87    /// or values will cause an error to be returned when connecting.
88    pub headers: Option<HashMap<String, String>>,
89    /// HTTP headers to include on every RPC call as binary gRPC metadata (encoded as base64).
90    ///
91    /// These must be valid binary gRPC metadata keys (and end with a `-bin` suffix). Invalid
92    /// header keys will cause an error to be returned when connecting.
93    pub binary_headers: Option<HashMap<String, Vec<u8>>>,
94    /// HTTP CONNECT proxy to use for this client.
95    pub http_connect_proxy: Option<HttpConnectProxyOptions>,
96    /// If set, DNS-based load balancing is enabled. When the target is a hostname (not an IP
97    /// literal), DNS is resolved to all addresses and requests are distributed across them.
98    /// Incompatible with `service_override` and `http_connect_proxy`. Setting either in addition
99    /// to this field is an error. Set to `None` to disable.
100    #[builder(required, default = Some(DnsLoadBalancingOptions::default()))]
101    pub dns_load_balancing: Option<DnsLoadBalancingOptions>,
102    /// If set true, error code labels will not be included on request failure metrics.
103    #[builder(default)]
104    pub disable_error_code_metric_tags: bool,
105    /// If set, all gRPC calls will be routed through the provided service.
106    pub service_override: Option<callback_based::CallbackBasedGrpcService>,
107    /// Controls transport-level gRPC compression for the client. Defaults to
108    /// [GrpcCompression::Gzip], which compresses outbound request bodies and accepts
109    /// compressed responses. Set to [GrpcCompression::None] to opt out.
110    /// If service_override is specified, is forced to `None`.
111    #[builder(default)]
112    pub grpc_compression: GrpcCompression,
113    /// Payload size limit options for this connection. Defaults to the standard warning thresholds;
114    /// disable an individual warning by setting its threshold to `0`.
115    /// NOTE: Experimental
116    #[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    // Internal / Core-based SDK only options below =============================================
128    /// If set true, get_system_info will not be called upon connection.
129    #[builder(default)]
130    #[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
131    pub(crate) skip_get_system_info: bool,
132    /// The name of the SDK being implemented on top of core. Is set as `client-name` header in
133    /// all RPC calls
134    #[builder(default = "temporal-rust".to_owned())]
135    #[cfg_attr(feature = "core-based-sdk", builder(setters(vis = "pub")))]
136    pub(crate) client_name: String,
137    // TODO [rust-sdk-branch]: SDK should set this to its version. Doing that probably easiest
138    // after adding proper client interceptors.
139    /// The version of the SDK being implemented on top of core. Is set as `client-version` header
140    /// in all RPC calls. The server decides if the client is supported based on this.
141    #[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// Bon does not propagate `doc(cfg)` to generated setters, so these docs-only methods forward to
147// renamed generated implementations.
148#[cfg(all(feature = "experimental", docsrs))]
149impl<S: connection_options_builder::State> ConnectionOptionsBuilder<S> {
150    /// Set the payload size limit options for this connection.
151    #[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    /// Set the payload size limit options for this connection from an optional value.
163    #[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// Setters/getters for fields that should only be touched by SDK implementers.
176#[cfg(feature = "core-based-sdk")]
177impl ConnectionOptions {
178    /// Set whether or not get_system_info will be called upon connection.
179    pub fn set_skip_get_system_info(&mut self, skip: bool) {
180        self.skip_get_system_info = skip;
181    }
182    /// Get whether or not get_system_info will be called upon connection.
183    pub fn get_skip_get_system_info(&self) -> bool {
184        self.skip_get_system_info
185    }
186    /// Get the name of the SDK being implemented on top of core.
187    pub fn get_client_name(&self) -> &str {
188        &self.client_name
189    }
190    /// Get the version of the SDK being implemented on top of core.
191    pub fn get_client_version(&self) -> &str {
192        &self.client_version
193    }
194}
195
196/// Options for [crate::Client::new].
197#[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    /// The namespace this client will be bound to.
202    #[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    /// The data converter used for serializing/deserializing payloads.
216    #[builder(default)]
217    pub data_converter: DataConverter,
218    /// Interceptors for high-level client operations, ordered outermost to innermost.
219    #[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    /// Register a type-erased client plugin.
227    ///
228    /// **Experimental:** This API may change or be removed.
229    pub fn plugin<P: Into<ErasedClientPlugin>>(mut self, plugin: P) -> Self {
230        self.plugins.push(plugin.into());
231        self
232    }
233
234    /// Register type-erased client plugins in iteration order.
235    ///
236    /// **Experimental:** This API may change or be removed.
237    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    /// Register a client-only plugin.
247    ///
248    /// **Experimental:** This API may change or be removed.
249    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    /// Return the registered plugins.
257    ///
258    /// This is intended for SDK integrations that propagate worker plugin registrations.
259    ///
260    /// **Experimental:** This API may change or be removed.
261    #[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/// Selects the transport-level compression used for gRPC calls. See
278/// [ConnectionOptions::grpc_compression].
279#[derive(Clone, Copy, Debug, PartialEq, Eq, Default)]
280#[non_exhaustive]
281pub enum GrpcCompression {
282    /// Do not compress requests or advertise acceptance of compressed responses.
283    None,
284    /// Gzip-compress outbound requests and accept gzip-compressed responses.
285    #[default]
286    Gzip,
287}
288
289/// Configuration options for TLS
290#[derive(Clone, bon::Builder)]
291#[non_exhaustive]
292pub struct TlsOptions {
293    /// Bytes representing the root CA certificate used by the server. If not set, and the server's
294    /// cert is issued by someone the operating system trusts, verification will still work (ex:
295    /// Cloud offering).
296    pub server_root_ca_cert: Option<Vec<u8>>,
297    /// Sets the domain name against which to verify the server's TLS certificate. If not provided,
298    /// the domain name will be extracted from the URL used to connect.
299    pub domain: Option<String>,
300    /// TLS info for the client. If specified, core will attempt to use mTLS.
301    ///
302    /// Mutually exclusive with [`client_cert_resolver`](TlsOptions::client_cert_resolver).
303    /// Setting both is an error.
304    pub client_tls_options: Option<ClientTlsOptions>,
305    /// Optional custom server certificate verifier. When set, this replaces the default
306    /// certificate verification and `server_root_ca_cert` is ignored.
307    ///
308    /// This is useful for:
309    /// - Certificate pinning
310    /// - Custom trust-domain validation (e.g., SAN-URI extraction)
311    /// - Federated root certificate stores
312    ///
313    /// # WARNING
314    /// Implementing a custom `ServerCertVerifier` can lead to severely insecure TLS connections
315    /// (e.g., disabling all validation or allowing man-in-the-middle attacks) if not done carefully.
316    /// Only use this if you know exactly what you are doing.
317    ///
318    /// The verifier must implement [`ServerCertVerifier`] from the `rustls` crate.
319    /// Note that `domain` is still respected for the `:authority` header / origin override
320    /// even when a custom verifier is set.
321    pub server_cert_verifier: Option<Arc<dyn ServerCertVerifier>>,
322    /// Optional dynamic client certificate resolver for transparent mTLS certificate rotation.
323    ///
324    /// Mutually exclusive with [`client_tls_options`](TlsOptions::client_tls_options).
325    /// Setting both is an error.
326    #[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/// If using mTLS, both the client cert and private key must be specified, this contains them.
362#[derive(Clone, bon::Builder)]
363#[non_exhaustive]
364pub struct ClientTlsOptions {
365    /// The certificate for this client, encoded as PEM
366    pub client_cert: Vec<u8>,
367    /// The private key for this client, encoded as PEM
368    pub client_private_key: Vec<u8>,
369}
370
371/// Client keep alive configuration.
372#[derive(Clone, Debug, PartialEq, bon::Builder)]
373#[non_exhaustive]
374pub struct ClientKeepAliveOptions {
375    /// Interval to send HTTP2 keep alive pings.
376    #[builder(default = Duration::from_secs(30))]
377    pub interval: Duration,
378    /// Timeout that the keep alive must be responded to within or the connection will be closed.
379    #[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/// Options for DNS-based load balancing.
390#[derive(Clone, Debug, PartialEq, bon::Builder)]
391#[non_exhaustive]
392pub struct DnsLoadBalancingOptions {
393    /// How often to re-resolve DNS. Defaults to 30 seconds.
394    #[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/// Payload size limit options for a connection.
405/// NOTE: Experimental
406#[cfg(feature = "experimental")]
407#[derive(Clone, Debug, PartialEq, bon::Builder)]
408#[non_exhaustive]
409pub struct PayloadLimitsOptions {
410    /// Warning threshold (bytes) for the size of an outbound payload-bearing field; over-threshold
411    /// fields are logged but still sent to server. Defaults to 512 KiB. Set to `0` to disable.
412    #[builder(default = DEFAULT_PAYLOADS_WARN_SIZE)]
413    pub payloads_warn_size: u64,
414    /// Warning threshold (bytes) for outbound memo sizes; over-threshold memos are logged but still
415    /// sent to server. Defaults to 2 KiB. Set to `0` to disable.
416    #[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    // Intentionally omit details here since they could leak a key if ever printed
429    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
430        write!(f, "ClientTlsOptions(..)")
431    }
432}
433
434/// Controls whether a closed workflow ID may be reused.
435#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
436#[non_exhaustive]
437pub enum WorkflowIdReusePolicy {
438    /// Use the server's default policy.
439    #[default]
440    Unspecified,
441    /// Allow starting a workflow using the same workflow ID.
442    AllowDuplicate,
443    /// Allow reuse only when the previous execution did not complete successfully.
444    AllowDuplicateFailedOnly,
445    /// Reject reuse of the workflow ID.
446    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/// Controls how starting a workflow resolves a conflict with a running workflow using the same
461/// workflow ID.
462#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
463#[non_exhaustive]
464pub enum WorkflowIdConflictPolicy {
465    /// Use the server's default policy.
466    #[default]
467    Unspecified,
468    /// Do not start a new workflow and return an already-started error.
469    Fail,
470    /// Do not start a new workflow and return a handle for the running workflow.
471    UseExisting,
472    /// Terminate the running workflow before starting a new one.
473    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/// Options for starting a workflow execution.
488#[derive(Debug, Clone, bon::Builder)]
489#[builder(start_fn = new, on(String, into))]
490#[non_exhaustive]
491pub struct WorkflowStartOptions {
492    /// The task queue to run the workflow on.
493    #[builder(start_fn)]
494    pub task_queue: String,
495
496    /// The workflow ID.
497    #[builder(start_fn)]
498    pub workflow_id: String,
499
500    /// Set the policy for reusing the workflow id
501    #[builder(default)]
502    pub id_reuse_policy: WorkflowIdReusePolicy,
503
504    /// Set the policy for how to resolve conflicts with running policies.
505    /// NOTE: This is ignored for child workflows.
506    #[builder(default)]
507    pub id_conflict_policy: WorkflowIdConflictPolicy,
508
509    /// Optionally set the execution timeout for the workflow
510    /// <https://docs.temporal.io/workflows/#workflow-execution-timeout>
511    pub execution_timeout: Option<Duration>,
512
513    /// Optionally indicates the default run timeout for a workflow run
514    pub run_timeout: Option<Duration>,
515
516    /// Optionally indicates the default task timeout for a workflow run
517    pub task_timeout: Option<Duration>,
518
519    /// Optionally set a cron schedule for the workflow
520    pub cron_schedule: Option<String>,
521
522    /// Additional search attributes for the workflow.
523    pub search_attributes: Option<SearchAttributes>,
524
525    /// Optionally enable Eager Workflow Start, a latency optimization using local workers.
526    #[builder(default)]
527    pub enable_eager_workflow_start: bool,
528
529    /// Optionally set a retry policy for the workflow
530    #[builder(into)]
531    pub retry_policy: Option<RetryPolicy>,
532
533    /// Links to associate with the workflow. Ex: References to a nexus operation.
534    #[builder(default)]
535    pub links: Vec<common::v1::Link>,
536
537    /// Callbacks that will be invoked upon workflow completion. For, ex, completing nexus
538    /// operations.
539    #[builder(default)]
540    pub completion_callbacks: Vec<common::v1::Callback>,
541
542    /// Priority for the workflow. Defaults to all-inherited (empty).
543    #[builder(default)]
544    pub priority: Priority,
545
546    /// Override the workflow's worker deployment routing. When unset, normal task queue
547    /// routing applies.
548    ///
549    /// **Experimental:** See [`VersioningOverride`] for the available routing behaviors.
550    pub versioning_override: Option<VersioningOverride>,
551
552    /// Headers to include with the start request.
553    pub header: Option<Header>,
554
555    /// Non-indexed values attached to the workflow, serialized with the client's data converter.
556    pub memo: Option<MemoValues>,
557
558    /// Single-line static summary for the workflow, shown in the Temporal UI.
559    pub static_summary: Option<String>,
560
561    /// Multi-line static details for the workflow, shown in the Temporal UI.
562    pub static_details: Option<String>,
563
564    /// Controls for the RPC used to start the workflow.
565    #[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/// Options for starting a workflow and sending it an update in one atomic operation.
623///
624/// See [crate::Client::start_update_with_start_workflow] and
625/// [crate::Client::execute_update_with_start_workflow].
626#[derive(Debug, Clone, bon::Builder)]
627#[builder(start_fn = new, on(String, into))]
628#[non_exhaustive]
629pub struct WorkflowUpdateWithStartOptions {
630    /// The task queue to run the workflow on.
631    #[builder(start_fn)]
632    pub task_queue: String,
633
634    /// The workflow ID.
635    #[builder(start_fn)]
636    pub workflow_id: String,
637
638    /// How to resolve a conflict with an already-running workflow. This is required so callers
639    /// explicitly choose whether an update may attach to an existing workflow.
640    #[builder(start_fn)]
641    pub id_conflict_policy: WorkflowIdConflictPolicy,
642
643    /// The policy for reusing the workflow ID after a workflow closes.
644    #[builder(default)]
645    pub id_reuse_policy: WorkflowIdReusePolicy,
646
647    /// The workflow execution timeout.
648    pub execution_timeout: Option<Duration>,
649
650    /// The workflow run timeout.
651    pub run_timeout: Option<Duration>,
652
653    /// The workflow task timeout.
654    pub task_timeout: Option<Duration>,
655
656    /// Search attributes for the workflow.
657    pub search_attributes: Option<SearchAttributes>,
658
659    /// The workflow retry policy.
660    #[builder(into)]
661    pub retry_policy: Option<RetryPolicy>,
662
663    /// Links to associate with the workflow.
664    #[builder(default)]
665    pub links: Vec<common::v1::Link>,
666
667    /// Callbacks invoked when the workflow completes.
668    #[builder(default)]
669    pub completion_callbacks: Vec<common::v1::Callback>,
670
671    /// Priority for the workflow. Defaults to all-inherited (empty).
672    #[builder(default)]
673    pub priority: Priority,
674
675    /// Override worker deployment routing when this operation starts a new workflow.
676    /// Does not change the routing of an already-running workflow.
677    ///
678    /// **Experimental:** See [`VersioningOverride`] for the available routing behaviors.
679    pub versioning_override: Option<VersioningOverride>,
680
681    /// Headers to include with the start operation.
682    pub start_header: Option<Header>,
683
684    /// Headers to include with the update operation.
685    pub update_header: Option<Header>,
686
687    /// Non-indexed values attached to the workflow, serialized with the client's data converter.
688    pub memo: Option<MemoValues>,
689
690    /// Single-line static summary for the workflow, shown in the Temporal UI.
691    pub static_summary: Option<String>,
692
693    /// Multi-line static details for the workflow, shown in the Temporal UI.
694    pub static_details: Option<String>,
695
696    /// Update ID for idempotency. If not provided, a UUID will be generated.
697    pub update_id: Option<String>,
698
699    /// Controls for the multi-operation RPC and, when executing the update, subsequent polling.
700    #[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/// Options for fetching workflow results
760#[derive(Debug, Clone, bon::Builder)]
761#[non_exhaustive]
762pub struct WorkflowGetResultOptions {
763    /// If true (the default), follows to the next workflow run in the execution chain while
764    /// retrieving results.
765    #[builder(default = true)]
766    pub follow_runs: bool,
767    /// Controls for each history RPC used to retrieve the result.
768    #[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/// Options for starting a workflow update.
781#[derive(Debug, Clone, Default, bon::Builder)]
782#[non_exhaustive]
783pub struct WorkflowExecuteUpdateOptions {
784    /// Update ID for idempotency.
785    pub update_id: Option<String>,
786    /// Headers to include.
787    pub header: Option<Header>,
788    /// Controls for the start-update and poll-update RPCs.
789    #[builder(default)]
790    pub rpc_options: RpcOptions,
791}
792
793/// Options for sending a signal to a workflow.
794#[derive(Debug, Clone, Default, bon::Builder)]
795#[non_exhaustive]
796pub struct WorkflowSignalOptions {
797    /// Request ID for idempotency. If not provided, a UUID will be generated.
798    pub request_id: Option<String>,
799    /// Headers to include with the signal.
800    pub header: Option<Header>,
801    /// Controls for the signal RPC.
802    #[builder(default)]
803    pub rpc_options: RpcOptions,
804}
805
806/// Controls when a workflow query should be rejected based on workflow state.
807#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
808#[non_exhaustive]
809pub enum QueryRejectCondition {
810    /// Use the server's default condition.
811    #[default]
812    Unspecified,
813    /// Do not reject the query based on workflow state.
814    None,
815    /// Reject the query if the workflow is not open.
816    NotOpen,
817    /// Reject the query if the workflow did not complete successfully.
818    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/// Options for querying a workflow.
833#[derive(Debug, Clone, Default, bon::Builder)]
834#[non_exhaustive]
835pub struct WorkflowQueryOptions {
836    /// Query reject condition. Determines when the query should be rejected
837    /// based on workflow state.
838    pub reject_condition: Option<QueryRejectCondition>,
839    /// Headers to include with the query.
840    pub header: Option<Header>,
841    /// Controls for the query RPC.
842    #[builder(default)]
843    pub rpc_options: RpcOptions,
844}
845
846/// Options for cancelling a workflow.
847#[derive(Debug, Clone, Default, bon::Builder)]
848#[builder(on(String, into))]
849#[non_exhaustive]
850pub struct WorkflowCancelOptions {
851    /// Reason for cancellation.
852    #[builder(default)]
853    pub reason: String,
854    /// Request ID for idempotency. If not provided, a UUID will be generated.
855    pub request_id: Option<String>,
856    /// Controls for the cancellation RPC.
857    #[builder(default)]
858    pub rpc_options: RpcOptions,
859}
860
861/// Options for terminating a workflow.
862#[derive(Debug, Clone, Default, bon::Builder)]
863#[builder(on(String, into))]
864#[non_exhaustive]
865pub struct WorkflowTerminateOptions {
866    /// Reason for termination.
867    #[builder(default)]
868    pub reason: String,
869    /// Additional details to include with the termination.
870    pub details: Option<Payloads>,
871    /// Controls for the termination RPC.
872    #[builder(default)]
873    pub rpc_options: RpcOptions,
874}
875
876/// Options for describing a workflow.
877#[derive(Debug, Clone, Default, bon::Builder)]
878#[non_exhaustive]
879pub struct WorkflowDescribeOptions {
880    /// Controls for the describe RPC.
881    #[builder(default)]
882    pub rpc_options: RpcOptions,
883}
884
885/// Default workflow execution retention for a Namespace is 3 days
886const DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD: Duration = Duration::from_secs(60 * 60 * 24 * 3);
887
888/// Controls whether archival is enabled for a namespace.
889#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
890#[non_exhaustive]
891pub enum ArchivalState {
892    /// Use the server's default archival state.
893    #[default]
894    Unspecified,
895    /// Disable archival.
896    Disabled,
897    /// Enable archival.
898    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/// Helper struct for `register_namespace`.
912#[derive(Clone, Debug, bon::Builder)]
913#[builder(on(String, into))]
914#[non_exhaustive]
915pub struct RegisterNamespaceOptions {
916    /// Name (required)
917    pub namespace: String,
918    /// Description (required)
919    pub description: String,
920    /// Owner's email
921    #[builder(default)]
922    pub owner_email: String,
923    /// Workflow execution retention period
924    #[builder(default = DEFAULT_WORKFLOW_EXECUTION_RETENTION_PERIOD)]
925    pub workflow_execution_retention_period: Duration,
926    /// Cluster settings
927    #[builder(default)]
928    pub clusters: Vec<ClusterReplicationConfig>,
929    /// Active cluster name
930    #[builder(default)]
931    pub active_cluster_name: String,
932    /// Custom Data
933    #[builder(default)]
934    pub data: HashMap<String, String>,
935    /// Security Token
936    #[builder(default)]
937    pub security_token: String,
938    /// Global namespace
939    #[builder(default)]
940    pub is_global_namespace: bool,
941    /// History Archival setting
942    #[builder(default = ArchivalState::Unspecified)]
943    pub history_archival_state: ArchivalState,
944    /// History Archival uri
945    #[builder(default)]
946    pub history_archival_uri: String,
947    /// Visibility Archival setting
948    #[builder(default = ArchivalState::Unspecified)]
949    pub visibility_archival_state: ArchivalState,
950    /// Visibility Archival uri
951    #[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/// Selects which workflow history events are returned when fetching history.
980#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
981#[non_exhaustive]
982pub enum HistoryEventFilterType {
983    /// Use the server's default filter.
984    #[default]
985    Unspecified,
986    /// Return all history events.
987    AllEvent,
988    /// Return only the workflow's close event.
989    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/// Options for fetching workflow history.
1003#[derive(Debug, Clone, Default, bon::Builder)]
1004#[non_exhaustive]
1005pub struct WorkflowFetchHistoryOptions {
1006    /// Whether to skip archival.
1007    #[builder(default)]
1008    pub skip_archival: bool,
1009    /// If set true, the fetch will wait for a new event before returning.
1010    #[builder(default)]
1011    pub wait_new_event: bool,
1012    /// Specifies which kind of events will be retrieved. Defaults to all events.
1013    #[builder(default = HistoryEventFilterType::AllEvent)]
1014    pub event_filter_type: HistoryEventFilterType,
1015    /// Controls for each history page RPC.
1016    #[builder(default)]
1017    pub rpc_options: RpcOptions,
1018}
1019
1020/// Options for starting an update without waiting for completion.
1021#[derive(Debug, Clone, Default, bon::Builder)]
1022#[non_exhaustive]
1023pub struct WorkflowStartUpdateOptions {
1024    /// Update ID for idempotency. If not provided, a UUID will be generated.
1025    pub update_id: Option<String>,
1026    /// Headers to include with the update.
1027    pub header: Option<Header>,
1028    /// Controls for the start-update RPC.
1029    #[builder(default)]
1030    pub rpc_options: RpcOptions,
1031}
1032
1033impl From<WorkflowExecuteUpdateOptions> for WorkflowStartUpdateOptions {
1034    /// Execute-update is start-update followed by waiting for the update result.
1035    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/// Options for listing workflows.
1045#[derive(Debug, Clone, Default, bon::Builder)]
1046#[non_exhaustive]
1047pub struct WorkflowListOptions {
1048    /// Maximum number of workflows to return.
1049    /// If not specified, returns all matching workflows.
1050    pub limit: Option<usize>,
1051    /// Controls for each list page RPC.
1052    #[builder(default)]
1053    pub rpc_options: RpcOptions,
1054}
1055
1056/// Options for counting workflows.
1057#[derive(Debug, Clone, Default, bon::Builder)]
1058#[non_exhaustive]
1059pub struct WorkflowCountOptions {
1060    /// Controls for the count RPC.
1061    #[builder(default)]
1062    pub rpc_options: RpcOptions,
1063}
1064
1065/// Options for starting a standalone activity.
1066#[derive(Clone, Debug, bon::Builder)]
1067#[builder(start_fn = new, on(String, into))]
1068#[non_exhaustive]
1069pub struct ActivityStartOptions {
1070    /// Task queue to run this activity on.
1071    #[builder(start_fn)]
1072    pub task_queue: String,
1073    /// Activity ID of the started activity. It's recommended to use a meaningful business ID.
1074    #[builder(start_fn)]
1075    pub id: String,
1076    /// Timeouts for activity completion.
1077    ///
1078    /// See [`ActivityCloseTimeouts`] for the meaning of each timeout variant.
1079    #[builder(start_fn)]
1080    pub close_timeouts: ActivityCloseTimeouts,
1081    /// If set, specifies maximum time the activity can wait in the task queue before being picked
1082    /// up by a worker. This timeout is non-retryable.
1083    pub schedule_to_start_timeout: Option<Duration>,
1084    /// If set, specifies maximum time between successful heartbeats.
1085    pub heartbeat_timeout: Option<Duration>,
1086    /// Controls how Activity is retried. If not set, the server will assign default retry policy.
1087    #[builder(into)]
1088    pub retry_policy: Option<RetryPolicy>,
1089    /// Priority to use when starting this activity.
1090    #[builder(default)]
1091    pub priority: Priority,
1092    /// Specifies behavior if there's a *closed* activity with the same ID.
1093    #[builder(default)]
1094    pub id_reuse_policy: ActivityIdReusePolicy,
1095    /// Specifies behavior if there's a *running* activity with the same ID. Note that there can
1096    /// only be one running activity for each Activity ID.
1097    #[builder(default)]
1098    pub id_conflict_policy: ActivityIdConflictPolicy,
1099    /// Search attributes for the activity.
1100    pub search_attributes: Option<SearchAttributes>,
1101    /// Headers to include with the start request.
1102    pub header: Option<Header>,
1103    /// Single-line static summary for the activity, shown in the Temporal UI.
1104    pub summary: Option<String>,
1105    /// Multi-line static details for the activity, shown in the Temporal UI.
1106    pub static_details: Option<String>,
1107    /// Time to wait before dispatching the first activity task.
1108    /// This delay is not applied to retry attempts.
1109    pub start_delay: Option<Duration>,
1110}
1111
1112impl ActivityStartOptions {
1113    /// Returns a builder with `close_timeouts` set to [`ActivityCloseTimeouts::StartToClose`].
1114    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    /// Returns a builder with `close_timeouts` set to [`ActivityCloseTimeouts::ScheduleToClose`].
1127    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/// Specifies behavior when starting a standalone activity if there's a *closed* activity with
1141/// the same ID. See [`ActivityStartOptions::id_reuse_policy`].
1142#[non_exhaustive]
1143#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
1144pub enum ActivityIdReusePolicy {
1145    #[default]
1146    /// Always allow starting an activity using the same activity ID. This is the default.
1147    AllowDuplicate,
1148    /// Allow starting an activity using the same ID only when the last execution did not complete
1149    /// successfully.
1150    AllowDuplicateFailedOnly,
1151    /// Do not permit re-use of the ID for this activity.
1152    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/// Specifies behavior when starting a standalone activity if there's a *running* activity with
1166/// the same ID. See [`ActivityStartOptions::id_conflict_policy`].
1167#[non_exhaustive]
1168#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
1169pub enum ActivityIdConflictPolicy {
1170    #[default]
1171    /// Don't start a new activity; instead return
1172    /// [`StartActivityError::AlreadyStarted`](crate::errors::StartActivityError::AlreadyStarted).
1173    Fail,
1174    /// Don't start a new activity; instead return a handle for the running activity.
1175    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/// Options for listing activities.
1188#[derive(Debug, Clone, Default, bon::Builder)]
1189#[non_exhaustive]
1190pub struct ActivityListOptions {}
1191
1192/// Options for counting activities.
1193#[derive(Debug, Clone, Default, bon::Builder)]
1194#[non_exhaustive]
1195pub struct ActivityCountOptions {}
1196
1197/// Controls which optional fields will be requested in
1198/// [`ActivityHandle::describe`](crate::ActivityHandle::describe) operation. The fields will be
1199/// present in returned [`ActivityExecutionDescription`](crate::ActivityExecutionDescription),
1200/// subject to data availability and server support.
1201///
1202/// Note that these fields contain payloads that can be arbitrarily large. It's recommended not to
1203/// include them unless they're needed.
1204#[derive(Debug, Clone, Default, bon::Builder)]
1205#[non_exhaustive]
1206pub struct ActivityDescribeOptions {
1207    /// If set and the activity received input, the input will be included.
1208    #[builder(default)]
1209    pub include_input: bool,
1210    /// If set and the activity is closed, the activity outcome will be included.
1211    #[builder(default)]
1212    pub include_outcome: bool,
1213    /// If set and the activity sent heartbeat details, the heartbeat details will be included.
1214    #[builder(default)]
1215    pub include_heartbeat_details: bool,
1216    /// If set and the activity has a failed attempt, the last failure will be included.
1217    #[builder(default)]
1218    pub include_last_failure: bool,
1219}
1220
1221/// Options for [`ActivityHandle::cancel`](crate::ActivityHandle::cancel).
1222#[derive(Debug, Clone, Default, bon::Builder)]
1223#[builder(on(String, into))]
1224#[non_exhaustive]
1225pub struct ActivityCancelOptions {
1226    /// Reason for cancellation. Can be empty.
1227    #[builder(default)]
1228    pub reason: String,
1229}
1230
1231/// Options for [`ActivityHandle::terminate`](crate::ActivityHandle::terminate).
1232#[derive(Debug, Clone, Default, bon::Builder)]
1233#[builder(on(String, into))]
1234#[non_exhaustive]
1235pub struct ActivityTerminateOptions {
1236    /// Reason for termination. Can be empty.
1237    #[builder(default)]
1238    pub reason: String,
1239}