Skip to main content

temporalio_client/
workflow_handle.rs

1use crate::{
2    CancelWorkflowInput, DescribeWorkflowInput, DescribeWorkflowOutput,
3    FetchWorkflowHistoryPageInput, FetchWorkflowHistoryPageOutput, HistoryEventFilterType,
4    NamespacedClient, Next, PollWorkflowUpdateInput, PollWorkflowUpdateOutput, QueryWorkflowInput,
5    QueryWorkflowOutput, RpcOptions, SignalWorkflowInput, StartWorkflowUpdateInput,
6    StartWorkflowUpdateOutput, TerminateWorkflowInput, WorkflowCancelOptions,
7    WorkflowDescribeOptions, WorkflowExecuteUpdateOptions, WorkflowExecutionStatus,
8    WorkflowFetchHistoryOptions, WorkflowGetResultOptions, WorkflowQueryOptions,
9    WorkflowSignalOptions, WorkflowStartUpdateOptions, WorkflowTerminateOptions,
10    errors::{
11        WorkflowGetResultError, WorkflowInteractionError, WorkflowQueryError, WorkflowUpdateError,
12    },
13    grpc::WorkflowService,
14    interceptors,
15};
16use futures_util::{TryStreamExt, future::BoxFuture, stream, stream::Stream};
17use std::{
18    collections::VecDeque,
19    fmt::Debug,
20    marker::PhantomData,
21    pin::Pin,
22    task::{Context, Poll},
23};
24pub use temporalio_common::UntypedWorkflow;
25use temporalio_common::{
26    HasWorkflowDefinition, QueryDefinition, SignalDefinition, UpdateDefinition, WorkflowDefinition,
27    data_converters::{
28        DataConverter, DecodablePayloads, GenericPayloadConverter, PayloadConversionError,
29        PayloadConverter, RawValue, SerializationContext, SerializationContextData,
30        WorkflowSerializationContext,
31    },
32    error::IncomingError,
33    payload_visitor::decode_payloads,
34    protos::{
35        coresdk::FromPayloadsExt,
36        proto_ts_to_system_time,
37        temporal::api::{
38            common::v1::{Header, Payload, Payloads, WorkflowExecution as ProtoWorkflowExecution},
39            enums::v1::{
40                HistoryEventFilterType as ProtoHistoryEventFilterType,
41                QueryRejectCondition as ProtoQueryRejectCondition,
42                UpdateWorkflowExecutionLifecycleStage,
43            },
44            history::{
45                self,
46                v1::{History, HistoryEvent, history_event::Attributes},
47            },
48            query::v1::WorkflowQuery,
49            sdk::v1::UserMetadata,
50            update::{self, v1::WaitPolicy},
51            workflow::v1 as workflow,
52            workflowservice::v1::{
53                DescribeWorkflowExecutionRequest, DescribeWorkflowExecutionResponse,
54                GetWorkflowExecutionHistoryRequest, PollWorkflowExecutionUpdateRequest,
55                QueryWorkflowRequest, RequestCancelWorkflowExecutionRequest,
56                SignalWorkflowExecutionRequest, TerminateWorkflowExecutionRequest,
57                UpdateWorkflowExecutionRequest,
58            },
59        },
60    },
61    search_attributes::SearchAttributes,
62};
63use tonic::IntoRequest;
64use uuid::Uuid;
65
66#[derive(Debug, Clone, Default, PartialEq, Eq)]
67struct DecodedUserMetadata {
68    summary: Option<String>,
69    details: Option<String>,
70}
71
72fn decode_user_metadata(
73    context: &SerializationContextData,
74    user_metadata: Option<UserMetadata>,
75) -> Result<DecodedUserMetadata, PayloadConversionError> {
76    let payload_converter = PayloadConverter::default();
77    let context = SerializationContext::new(context, &payload_converter);
78    let (summary, details) = user_metadata
79        .map(|metadata| (metadata.summary, metadata.details))
80        .unwrap_or_default();
81    Ok(DecodedUserMetadata {
82        summary: match summary {
83            Some(payload) => Some(payload_converter.from_payload(&context, payload)?),
84            None => None,
85        },
86        details: match details {
87            Some(payload) => Some(payload_converter.from_payload(&context, payload)?),
88            None => None,
89        },
90    })
91}
92
93/// Details attached to a cancelled or terminated workflow result.
94#[derive(Clone, Debug)]
95#[non_exhaustive]
96pub struct WorkflowResultDetails {
97    payloads: DecodablePayloads,
98}
99
100impl WorkflowResultDetails {
101    async fn new(
102        payloads: Vec<Payload>,
103        data_converter: &DataConverter,
104    ) -> Result<Self, PayloadConversionError> {
105        let payloads = data_converter
106            .codec()
107            .decode(
108                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
109                payloads,
110            )
111            .await?;
112        Ok(Self {
113            payloads: DecodablePayloads::new(
114                payloads,
115                data_converter.payload_converter().clone(),
116                SerializationContextData::Workflow(WorkflowSerializationContext::new()),
117            ),
118        })
119    }
120
121    /// Deserialize the details into a typed value using the client's payload converter.
122    pub fn deserialize<T: temporalio_common::data_converters::TemporalDeserializable + 'static>(
123        &self,
124    ) -> Result<T, PayloadConversionError> {
125        self.payloads.deserialize()
126    }
127
128    /// Returns the codec-decoded payloads.
129    pub fn raw(&self) -> &[Payload] {
130        self.payloads.raw()
131    }
132
133    /// Consume these details and return their codec-decoded payloads.
134    pub fn into_raw(self) -> RawValue {
135        self.payloads.into_raw()
136    }
137}
138
139/// Enumerates terminal states for a particular workflow execution
140#[derive(Debug)]
141#[allow(clippy::large_enum_variant)]
142pub enum WorkflowExecutionResult<T> {
143    /// The workflow finished successfully
144    Succeeded(T),
145    /// The workflow finished in failure
146    Failed(IncomingError),
147    /// The workflow was cancelled
148    Cancelled {
149        /// Details provided at cancellation time
150        details: WorkflowResultDetails,
151    },
152    /// The workflow was terminated
153    Terminated {
154        /// Details provided at termination time
155        details: WorkflowResultDetails,
156    },
157    /// The workflow timed out
158    TimedOut,
159    /// The workflow continued as new
160    ContinuedAsNew,
161}
162
163/// Description of a workflow execution returned by `WorkflowHandle::describe`.
164///
165/// Access to the underlying Protobuf message is provided by [`raw`](Self::raw).
166#[derive(Debug, Clone)]
167pub struct WorkflowExecutionDescription {
168    /// The raw proto response from the server.
169    pub raw_description: DescribeWorkflowExecutionResponse,
170    history_length: usize,
171    static_summary: Option<String>,
172    static_details: Option<String>,
173    data_converter: DataConverter,
174}
175
176impl WorkflowExecutionDescription {
177    async fn new(
178        mut raw_description: DescribeWorkflowExecutionResponse,
179        data_converter: &DataConverter,
180    ) -> Result<Self, PayloadConversionError> {
181        let raw_user_metadata = raw_description
182            .execution_config
183            .as_ref()
184            .and_then(|cfg| cfg.user_metadata.clone());
185        decode_payloads(
186            &mut raw_description,
187            data_converter.codec(),
188            &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
189        )
190        .await?;
191        let decoded_metadata = decode_user_metadata(
192            &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
193            raw_user_metadata,
194        )?;
195        let history_length_raw = raw_description
196            .workflow_execution_info
197            .as_ref()
198            .map(|info| info.history_length)
199            .unwrap_or(0);
200        let history_length = history_length_raw.try_into().map_err(|_| {
201            PayloadConversionError::EncodingError(
202                format!("workflow history_length must be non-negative, got {history_length_raw}")
203                    .into(),
204            )
205        })?;
206        Ok(Self {
207            raw_description,
208            history_length,
209            static_summary: decoded_metadata.summary,
210            static_details: decoded_metadata.details,
211            data_converter: data_converter.clone(),
212        })
213    }
214
215    /// The workflow ID.
216    pub fn id(&self) -> &str {
217        self.execution().workflow_id.as_str()
218    }
219
220    /// The run ID.
221    pub fn run_id(&self) -> &str {
222        self.execution().run_id.as_str()
223    }
224
225    /// The workflow type name.
226    pub fn workflow_type(&self) -> &str {
227        self.workflow_type_info().name.as_str()
228    }
229
230    /// The current status of the workflow execution.
231    pub fn status(&self) -> WorkflowExecutionStatus {
232        WorkflowExecutionStatus::from_raw(self.workflow_info().status)
233    }
234
235    /// When the workflow was created.
236    pub fn start_time(&self) -> Option<std::time::SystemTime> {
237        self.workflow_info()
238            .start_time
239            .as_ref()
240            .and_then(proto_ts_to_system_time)
241    }
242
243    /// When the workflow run started or should start.
244    pub fn execution_time(&self) -> Option<std::time::SystemTime> {
245        self.workflow_info()
246            .execution_time
247            .as_ref()
248            .and_then(proto_ts_to_system_time)
249    }
250
251    /// When the workflow was closed, if closed.
252    pub fn close_time(&self) -> Option<std::time::SystemTime> {
253        self.workflow_info()
254            .close_time
255            .as_ref()
256            .and_then(proto_ts_to_system_time)
257    }
258
259    /// The task queue the workflow runs on.
260    pub fn task_queue(&self) -> &str {
261        self.workflow_info().task_queue.as_str()
262    }
263
264    /// Number of events in history.
265    pub fn history_length(&self) -> usize {
266        self.history_length
267    }
268
269    /// Workflow memo decoded with the client's payload converter.
270    pub fn memo(&self) -> crate::Memo {
271        crate::Memo::from_raw(
272            self.workflow_info().memo.clone(),
273            self.data_converter.payload_converter().clone(),
274            SerializationContextData::Workflow(WorkflowSerializationContext::new()),
275        )
276    }
277
278    /// Parent workflow ID, if this is a child workflow.
279    pub fn parent_id(&self) -> Option<&str> {
280        self.workflow_info()
281            .parent_execution
282            .as_ref()
283            .map(|e| e.workflow_id.as_str())
284    }
285
286    /// Parent run ID, if this is a child workflow.
287    pub fn parent_run_id(&self) -> Option<&str> {
288        self.workflow_info()
289            .parent_execution
290            .as_ref()
291            .map(|e| e.run_id.as_str())
292    }
293
294    /// Search attributes on the workflow.
295    pub fn search_attributes(&self) -> SearchAttributes {
296        self.workflow_info()
297            .search_attributes
298            .as_ref()
299            .map(SearchAttributes::from_proto)
300            .unwrap_or_default()
301    }
302
303    /// Static summary configured on the workflow, if present.
304    pub fn static_summary(&self) -> Option<&str> {
305        self.static_summary.as_deref()
306    }
307
308    /// Static details configured on the workflow, if present.
309    pub fn static_details(&self) -> Option<&str> {
310        self.static_details.as_deref()
311    }
312
313    /// Access the raw proto for additional fields not exposed via accessors.
314    pub fn raw(&self) -> &DescribeWorkflowExecutionResponse {
315        &self.raw_description
316    }
317
318    /// Consume the wrapper and return the raw proto.
319    pub fn into_raw(self) -> DescribeWorkflowExecutionResponse {
320        self.raw_description
321    }
322
323    fn workflow_info(&self) -> &workflow::WorkflowExecutionInfo {
324        self.raw_description
325            .workflow_execution_info
326            .as_ref()
327            .expect("describe response missing workflow_execution_info")
328    }
329
330    fn execution(&self) -> &ProtoWorkflowExecution {
331        self.workflow_info()
332            .execution
333            .as_ref()
334            .expect("describe response missing workflow_execution_info.execution")
335    }
336
337    fn workflow_type_info(
338        &self,
339    ) -> &temporalio_common::protos::temporal::api::common::v1::WorkflowType {
340        self.workflow_info()
341            .r#type
342            .as_ref()
343            .expect("describe response missing workflow_execution_info.type")
344    }
345}
346
347/// Workflow execution history returned by [`WorkflowHandle::fetch_history`].
348///
349/// Events and their containing pages are fetched lazily as this stream is polled. Use
350/// [`into_events`](Self::into_events) to fetch and collect all events at once.
351#[derive(derive_more::Debug)]
352pub struct WorkflowHistory {
353    #[debug(skip)]
354    inner: Pin<Box<dyn Stream<Item = Result<HistoryEvent, WorkflowInteractionError>> + Send>>,
355    workflow_id: Option<String>,
356}
357
358impl From<history::v1::History> for WorkflowHistory {
359    fn from(history: history::v1::History) -> Self {
360        let workflow_id =
361            history
362                .events
363                .first()
364                .and_then(|event| match event.attributes.as_ref() {
365                    Some(Attributes::WorkflowExecutionStartedEventAttributes(attributes))
366                        if !attributes.workflow_id.is_empty() =>
367                    {
368                        Some(attributes.workflow_id.clone())
369                    }
370                    _ => None,
371                });
372        Self {
373            inner: Box::pin(stream::iter(history.events.into_iter().map(Ok))),
374            workflow_id,
375        }
376    }
377}
378
379impl Stream for WorkflowHistory {
380    type Item = Result<HistoryEvent, WorkflowInteractionError>;
381
382    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
383        self.inner.as_mut().poll_next(cx)
384    }
385}
386
387/// Error fetching or converting a workflow history.
388#[derive(Debug, thiserror::Error)]
389#[non_exhaustive]
390pub enum WorkflowHistoryError {
391    /// Fetching the workflow history failed.
392    #[error("failed to fetch workflow history: {0}")]
393    Fetch(#[from] WorkflowInteractionError),
394    /// Converting the workflow history JSON failed.
395    #[error("failed to convert workflow history JSON: {0}")]
396    Json(#[from] serde_json::Error),
397}
398
399impl WorkflowHistory {
400    /// Decode a workflow history from JSON bytes.
401    pub fn from_json(bytes: &[u8]) -> Result<Self, WorkflowHistoryError> {
402        let history: History = serde_json::from_slice(bytes)?;
403        Ok(history.into())
404    }
405
406    /// Fetch all remaining events and encode this workflow history as JSON bytes.
407    pub async fn to_json(self) -> Result<Vec<u8>, WorkflowHistoryError> {
408        Ok(serde_json::to_vec(&History {
409            events: self.into_events().await?,
410        })?)
411    }
412
413    /// Return the workflow ID when it is known.
414    pub fn workflow_id(&self) -> Option<&str> {
415        self.workflow_id.as_deref()
416    }
417
418    /// Fetch all remaining history pages and collect their events.
419    pub async fn into_events(self) -> Result<Vec<HistoryEvent>, WorkflowInteractionError> {
420        self.inner.try_collect().await
421    }
422}
423
424/// A workflow handle which can refer to a specific workflow run, or a chain of workflow runs with
425/// the same workflow id.
426#[derive(Clone)]
427pub struct WorkflowHandle<ClientT, W> {
428    client: ClientT,
429    info: WorkflowExecutionInfo,
430
431    _wf_type: PhantomData<W>,
432}
433
434impl<CT, W> WorkflowHandle<CT, W> {
435    /// Return the run id of the Workflow Execution pointed at by this handle, if there is one.
436    pub fn run_id(&self) -> Option<&str> {
437        self.info.run_id.as_deref()
438    }
439}
440
441/// Holds needed information to refer to a specific workflow run, or workflow execution chain
442#[derive(Debug, Clone, bon::Builder)]
443#[builder(on(String, into), state_mod(vis = "pub"))]
444#[non_exhaustive]
445pub struct WorkflowExecutionInfo {
446    /// Namespace the workflow lives in.
447    pub namespace: String,
448    /// The workflow's id.
449    pub workflow_id: String,
450    /// If set, target this specific run of the workflow.
451    pub run_id: Option<String>,
452    /// Run ID used for cancellation and termination to ensure they happen on a workflow starting
453    /// with this run ID. This can be set when getting a workflow handle. When starting a workflow,
454    /// this is set as the resulting run ID if no start signal was provided.
455    pub first_execution_run_id: Option<String>,
456}
457
458impl WorkflowExecutionInfo {
459    /// Bind the workflow info to a specific client, turning it into a workflow handle
460    pub fn bind_untyped<CT>(self, client: CT) -> UntypedWorkflowHandle<CT>
461    where
462        CT: WorkflowService + Clone,
463    {
464        UntypedWorkflowHandle::new(client, self)
465    }
466}
467
468/// A workflow handle to a workflow with unknown types. Uses single argument raw payloads for input
469/// and output.
470pub type UntypedWorkflowHandle<CT> = WorkflowHandle<CT, UntypedWorkflow>;
471
472/// Marker type for sending untyped signals. Stores the signal name for runtime lookup.
473///
474/// Use with `handle.signal(UntypedSignal::new("signal_name"), raw_payload)`.
475pub struct UntypedSignal<W> {
476    name: String,
477    _wf: PhantomData<W>,
478}
479
480impl<W> UntypedSignal<W> {
481    /// Create a new `UntypedSignal` with the given signal name.
482    pub fn new(name: impl Into<String>) -> Self {
483        Self {
484            name: name.into(),
485            _wf: PhantomData,
486        }
487    }
488}
489
490impl<W: WorkflowDefinition> SignalDefinition for UntypedSignal<W> {
491    type Workflow = W;
492    type Input = RawValue;
493
494    fn name(&self) -> &str {
495        &self.name
496    }
497}
498
499/// Marker type for sending untyped queries. Stores the query name for runtime lookup.
500///
501/// Use with `handle.query(UntypedQuery::new("query_name"), raw_payload)`.
502pub struct UntypedQuery<W> {
503    name: String,
504    _wf: PhantomData<W>,
505}
506
507impl<W> UntypedQuery<W> {
508    /// Create a new `UntypedQuery` with the given query name.
509    pub fn new(name: impl Into<String>) -> Self {
510        Self {
511            name: name.into(),
512            _wf: PhantomData,
513        }
514    }
515}
516
517impl<W: WorkflowDefinition> QueryDefinition for UntypedQuery<W> {
518    type Workflow = W;
519    type Input = RawValue;
520    type Output = RawValue;
521
522    fn name(&self) -> &str {
523        &self.name
524    }
525}
526
527/// Marker type for sending untyped updates. Stores the update name for runtime lookup.
528///
529/// Use with `handle.update(UntypedUpdate::new("update_name"), raw_payload)`.
530pub struct UntypedUpdate<W> {
531    name: String,
532    _wf: PhantomData<W>,
533}
534
535impl<W> UntypedUpdate<W> {
536    /// Create a new `UntypedUpdate` with the given update name.
537    pub fn new(name: impl Into<String>) -> Self {
538        Self {
539            name: name.into(),
540            _wf: PhantomData,
541        }
542    }
543}
544
545impl<W: WorkflowDefinition> UpdateDefinition for UntypedUpdate<W> {
546    type Workflow = W;
547    type Input = RawValue;
548    type Output = RawValue;
549
550    fn name(&self) -> &str {
551        &self.name
552    }
553}
554
555/// Shared by [WorkflowHandle::start_update] and the client's update-with-start, which sends the
556/// same update request as one of its operations. Update starts always wait for the update to be
557/// accepted; results are waited on separately via the update handle.
558#[allow(clippy::too_many_arguments)]
559pub(crate) fn build_update_workflow_request(
560    namespace: String,
561    identity: String,
562    workflow_id: String,
563    run_id: String,
564    update_id: String,
565    update_name: String,
566    header: Option<Header>,
567    payloads: Vec<Payload>,
568) -> UpdateWorkflowExecutionRequest {
569    UpdateWorkflowExecutionRequest {
570        namespace,
571        workflow_execution: Some(ProtoWorkflowExecution {
572            workflow_id,
573            run_id,
574        }),
575        wait_policy: Some(WaitPolicy {
576            lifecycle_stage: UpdateWorkflowExecutionLifecycleStage::Accepted.into(),
577        }),
578        request: Some(update::v1::Request {
579            meta: Some(update::v1::Meta {
580                update_id,
581                identity,
582            }),
583            input: Some(update::v1::Input {
584                header,
585                name: update_name,
586                args: Some(Payloads { payloads }),
587            }),
588            ..Default::default()
589        }),
590        ..Default::default()
591    }
592}
593
594impl<CT, W> WorkflowHandle<CT, W>
595where
596    CT: WorkflowService + Clone,
597    W: HasWorkflowDefinition,
598{
599    /// Create a workflow handle from a client and identifying information.
600    pub fn new(client: CT, info: WorkflowExecutionInfo) -> Self {
601        Self {
602            client,
603            info,
604            _wf_type: PhantomData::<W>,
605        }
606    }
607
608    /// Get the workflow execution info
609    pub fn info(&self) -> &WorkflowExecutionInfo {
610        &self.info
611    }
612
613    /// Get the client attached to this handle
614    pub fn client(&self) -> &CT {
615        &self.client
616    }
617
618    /// Await the result of the workflow execution
619    pub async fn get_result(
620        &self,
621        opts: WorkflowGetResultOptions,
622    ) -> Result<W::Output, WorkflowGetResultError>
623    where
624        CT: WorkflowService + NamespacedClient + Clone + 'static,
625    {
626        let raw = self.get_result_raw(opts).await?;
627        match raw {
628            WorkflowExecutionResult::Succeeded(v) => Ok(v),
629            WorkflowExecutionResult::Failed(f) => Err(WorkflowGetResultError::Failed(Box::new(f))),
630            WorkflowExecutionResult::Cancelled { details } => {
631                Err(WorkflowGetResultError::Cancelled { details })
632            }
633            WorkflowExecutionResult::Terminated { details } => {
634                Err(WorkflowGetResultError::Terminated { details })
635            }
636            WorkflowExecutionResult::TimedOut => Err(WorkflowGetResultError::TimedOut),
637            WorkflowExecutionResult::ContinuedAsNew => Err(WorkflowGetResultError::ContinuedAsNew),
638        }
639    }
640
641    /// Await the result of the workflow execution, returning the full
642    /// [`WorkflowExecutionResult`] enum for callers that need to inspect non-success outcomes
643    /// directly.
644    async fn get_result_raw(
645        &self,
646        opts: WorkflowGetResultOptions,
647    ) -> Result<WorkflowExecutionResult<W::Output>, WorkflowInteractionError>
648    where
649        CT: WorkflowService + NamespacedClient + Clone + 'static,
650    {
651        let mut run_id = self.info.run_id.clone().unwrap_or_default();
652        let fetch_opts = WorkflowFetchHistoryOptions::builder()
653            .skip_archival(true)
654            .wait_new_event(true)
655            .event_filter_type(HistoryEventFilterType::CloseEvent)
656            .rpc_options(opts.rpc_options.clone())
657            .build();
658
659        loop {
660            let history = self.fetch_history_for_run(&run_id, fetch_opts.clone());
661            let mut events = history.into_events().await?;
662
663            if events.is_empty() {
664                continue;
665            }
666
667            let event_attrs = events.pop().and_then(|ev| ev.attributes);
668
669            macro_rules! follow {
670                ($attrs:ident) => {
671                    if opts.follow_runs && $attrs.new_execution_run_id != "" {
672                        run_id = $attrs.new_execution_run_id;
673                        continue;
674                    }
675                };
676            }
677
678            let dc = self.client.data_converter();
679
680            break match event_attrs {
681                Some(Attributes::WorkflowExecutionCompletedEventAttributes(attrs)) => {
682                    follow!(attrs);
683                    let payload = attrs
684                        .result
685                        .and_then(|p| p.payloads.into_iter().next())
686                        .unwrap_or_default();
687                    let result: W::Output = dc
688                        .from_payload(&SerializationContextData::Workflow(WorkflowSerializationContext::new()), payload)
689                        .await?;
690                    Ok(WorkflowExecutionResult::Succeeded(result))
691                }
692                Some(Attributes::WorkflowExecutionFailedEventAttributes(attrs)) => {
693                    follow!(attrs);
694                    let mut failure = attrs.failure.unwrap_or_default();
695                    decode_payloads(
696                        &mut failure,
697                        dc.codec(),
698                        &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
699                    )
700                    .await?;
701                    let error = dc.failure_converter().to_error(
702                        failure,
703                        dc.payload_converter(),
704                        &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
705                    )?;
706                    Ok(WorkflowExecutionResult::Failed(error))
707                }
708                Some(Attributes::WorkflowExecutionCanceledEventAttributes(attrs)) => {
709                    Ok(WorkflowExecutionResult::Cancelled {
710                        details: WorkflowResultDetails::new(Vec::from_payloads(attrs.details), dc)
711                            .await?,
712                    })
713                }
714                Some(Attributes::WorkflowExecutionTimedOutEventAttributes(attrs)) => {
715                    follow!(attrs);
716                    Ok(WorkflowExecutionResult::TimedOut)
717                }
718                Some(Attributes::WorkflowExecutionTerminatedEventAttributes(attrs)) => {
719                    Ok(WorkflowExecutionResult::Terminated {
720                        details: WorkflowResultDetails::new(Vec::from_payloads(attrs.details), dc)
721                            .await?,
722                    })
723                }
724                Some(Attributes::WorkflowExecutionContinuedAsNewEventAttributes(attrs)) => {
725                    if opts.follow_runs {
726                        if !attrs.new_execution_run_id.is_empty() {
727                            run_id = attrs.new_execution_run_id;
728                            continue;
729                        } else {
730                            return Err(WorkflowInteractionError::Other(
731                                "New execution run id was empty in continue as new event!".into(),
732                            ));
733                        }
734                    } else {
735                        Ok(WorkflowExecutionResult::ContinuedAsNew)
736                    }
737                }
738                o => Err(WorkflowInteractionError::Other(
739                    format!(
740                        "Server returned an event that didn't match the CloseEvent filter. \
741                         This is either a server bug or a new event the SDK does not understand. \
742                         Event details: {o:?}"
743                    )
744                    .into(),
745                )),
746            };
747        }
748    }
749
750    /// Send a signal to the workflow
751    pub async fn signal<S>(
752        &self,
753        signal: S,
754        input: S::Input,
755        opts: WorkflowSignalOptions,
756    ) -> Result<(), WorkflowInteractionError>
757    where
758        CT: WorkflowService + NamespacedClient + Clone,
759        S: SignalDefinition<Workflow = W::Run>,
760        S::Input: Send,
761    {
762        interceptors::call_signal_workflow(
763            self.client.client_interceptors(),
764            SignalWorkflowInput::new(
765                self.info.workflow_id.clone(),
766                self.info.run_id.clone().unwrap_or_default(),
767                signal.name().to_string(),
768                input,
769                opts,
770            ),
771            Next::new({
772                let mut client = self.client.clone();
773                move |input: SignalWorkflowInput| -> BoxFuture<
774                    '_,
775                    Result<(), WorkflowInteractionError>,
776                > {
777                    Box::pin(async move {
778                        let (workflow_id, run_id, signal_name, args, options) =
779                            input.into_parts();
780                        let data_converter = client.data_converter().clone();
781                        let unencoded_payloads = {
782                            let payload_converter = data_converter.payload_converter();
783                            let context_data = SerializationContextData::Workflow(
784                                WorkflowSerializationContext::new(),
785                            );
786                            let context =
787                                SerializationContext::new(&context_data, payload_converter);
788                            args.serialize_payloads(&context)
789                        };
790                        drop(args);
791                        let payloads = data_converter
792                            .codec()
793                            .encode(&SerializationContextData::Workflow(WorkflowSerializationContext::new()), unencoded_payloads?)
794                            .await?;
795                        let mut request = SignalWorkflowExecutionRequest {
796                            namespace: client.namespace(),
797                            workflow_execution: Some(ProtoWorkflowExecution {
798                                workflow_id,
799                                run_id,
800                            }),
801                            signal_name,
802                            input: Some(Payloads { payloads }),
803                            identity: client.identity(),
804                            request_id: options
805                                .request_id
806                                .unwrap_or_else(|| Uuid::new_v4().to_string()),
807                            header: options.header,
808                            ..Default::default()
809                        }
810                        .into_request();
811                        options.rpc_options.apply_to(&mut request);
812                        WorkflowService::signal_workflow_execution(&mut client, request)
813                            .await
814                            .map_err(WorkflowInteractionError::from_status)?;
815                        Ok(())
816                    })
817                }
818            }),
819        )
820        .await
821    }
822
823    /// Query the workflow
824    pub async fn query<Q>(
825        &self,
826        query: Q,
827        input: Q::Input,
828        opts: WorkflowQueryOptions,
829    ) -> Result<Q::Output, WorkflowQueryError>
830    where
831        CT: WorkflowService + NamespacedClient + Clone,
832        Q: QueryDefinition<Workflow = W::Run>,
833        Q::Input: Send,
834    {
835        let output = interceptors::call_query_workflow(
836            self.client.client_interceptors(),
837            QueryWorkflowInput::new(
838                self.info.workflow_id.clone(),
839                self.info.run_id.clone().unwrap_or_default(),
840                query.name().to_string(),
841                input,
842                opts,
843            ),
844            Next::new({
845                let mut client = self.client.clone();
846                move |input: QueryWorkflowInput| -> BoxFuture<
847                    '_,
848                    Result<QueryWorkflowOutput, WorkflowQueryError>,
849                > {
850                    Box::pin(async move {
851                        let (workflow_id, run_id, query_name, args, options) = input.into_parts();
852                        let data_converter = client.data_converter().clone();
853                        let unencoded_payloads = {
854                            let payload_converter = data_converter.payload_converter();
855                            let context_data = SerializationContextData::Workflow(
856                                WorkflowSerializationContext::new(),
857                            );
858                            let context =
859                                SerializationContext::new(&context_data, payload_converter);
860                            args.serialize_payloads(&context)
861                        };
862                        drop(args);
863                        let payloads = data_converter
864                            .codec()
865                            .encode(&SerializationContextData::Workflow(WorkflowSerializationContext::new()), unencoded_payloads?)
866                            .await?;
867                        let mut request = QueryWorkflowRequest {
868                            namespace: client.namespace(),
869                            execution: Some(ProtoWorkflowExecution {
870                                workflow_id,
871                                run_id,
872                            }),
873                            query: Some(WorkflowQuery {
874                                query_type: query_name,
875                                query_args: Some(Payloads { payloads }),
876                                header: options.header,
877                            }),
878                            query_reject_condition: options
879                                .reject_condition
880                                .map(|condition| ProtoQueryRejectCondition::from(condition) as i32)
881                                .unwrap_or(ProtoQueryRejectCondition::None as i32),
882                        }
883                        .into_request();
884                        options.rpc_options.apply_to(&mut request);
885                        let response = client
886                            .query_workflow(request)
887                            .await
888                            .map_err(WorkflowQueryError::from_status)?
889                            .into_inner();
890                        Ok(QueryWorkflowOutput::new(response))
891                    })
892                }
893            }),
894        )
895        .await?;
896        let response = output.response;
897
898        if let Some(rejected) = response.query_rejected {
899            return Err(WorkflowQueryError::Rejected {
900                status: (rejected.status != 0)
901                    .then(|| WorkflowExecutionStatus::from_raw(rejected.status)),
902            });
903        }
904
905        let result_payloads = response
906            .query_result
907            .map(|p| p.payloads)
908            .unwrap_or_default();
909
910        self.client
911            .data_converter()
912            .from_payloads(
913                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
914                result_payloads,
915            )
916            .await
917            .map_err(WorkflowQueryError::from)
918    }
919
920    /// Send an update to the workflow and wait for it to complete, returning the result.
921    pub async fn execute_update<U>(
922        &self,
923        update: U,
924        input: U::Input,
925        options: WorkflowExecuteUpdateOptions,
926    ) -> Result<U::Output, WorkflowUpdateError>
927    where
928        CT: WorkflowService + NamespacedClient + Clone,
929        U: UpdateDefinition<Workflow = W::Run>,
930        U::Input: Send,
931        U::Output: 'static,
932    {
933        let rpc_options = options.rpc_options.clone();
934        let handle = self.start_update(update, input, options.into()).await?;
935        handle.get_result(rpc_options).await
936    }
937
938    /// Start an update and return a handle without waiting for completion.
939    /// Use `execute_update()` if you want to wait for the result immediately.
940    pub async fn start_update<U>(
941        &self,
942        update: U,
943        input: U::Input,
944        options: WorkflowStartUpdateOptions,
945    ) -> Result<WorkflowUpdateHandle<CT, U::Output>, WorkflowUpdateError>
946    where
947        CT: WorkflowService + NamespacedClient + Clone,
948        U: UpdateDefinition<Workflow = W::Run>,
949        U::Input: Send,
950    {
951        let output = interceptors::call_start_workflow_update(
952            self.client.client_interceptors(),
953            StartWorkflowUpdateInput::new(
954                self.info().workflow_id.clone(),
955                self.info().run_id.clone().unwrap_or_default(),
956                update.name().to_string(),
957                input,
958                options,
959            ),
960            Next::new({
961                let mut client = self.client.clone();
962                move |input: StartWorkflowUpdateInput| -> BoxFuture<
963                        '_,
964                        Result<StartWorkflowUpdateOutput, WorkflowUpdateError>,
965                    > {
966                        Box::pin(async move {
967                            let (workflow_id, run_id, update_name, args, options) =
968                                input.into_parts();
969                            let data_converter = client.data_converter().clone();
970                            let unencoded_payloads = {
971                                let payload_converter = data_converter.payload_converter();
972                                let context_data = SerializationContextData::Workflow(
973                                    WorkflowSerializationContext::new(),
974                                );
975                                let context =
976                                    SerializationContext::new(&context_data, payload_converter);
977                                args.serialize_payloads(&context)
978                            };
979                            drop(args);
980                            let payloads = data_converter
981                                .codec()
982                                .encode(
983                                    &SerializationContextData::Workflow(
984                                        WorkflowSerializationContext::new(),
985                                    ),
986                                    unencoded_payloads?,
987                                )
988                                .await?;
989                            let update_id = options
990                                .update_id
991                                .unwrap_or_else(|| Uuid::new_v4().to_string());
992                            let request = build_update_workflow_request(
993                                client.namespace(),
994                                client.identity(),
995                                workflow_id.clone(),
996                                run_id,
997                                update_id.clone(),
998                                update_name,
999                                options.header,
1000                                payloads,
1001                            );
1002                            let response = loop {
1003                                let mut rpc_request = request.clone().into_request();
1004                                options.rpc_options.apply_to(&mut rpc_request);
1005                                let response = WorkflowService::update_workflow_execution(
1006                                    &mut client,
1007                                    rpc_request,
1008                                )
1009                                .await
1010                                .map_err(WorkflowUpdateError::from_status)?
1011                                .into_inner();
1012                                if response.stage
1013                                    >= UpdateWorkflowExecutionLifecycleStage::Accepted as i32
1014                                {
1015                                    break response;
1016                                }
1017                            };
1018                            let run_id = response
1019                                .update_ref
1020                                .as_ref()
1021                                .and_then(|reference| reference.workflow_execution.as_ref())
1022                                .map(|execution| execution.run_id.clone())
1023                                .filter(|run_id| !run_id.is_empty());
1024                            Ok(StartWorkflowUpdateOutput::new(
1025                                update_id,
1026                                workflow_id,
1027                                run_id,
1028                                response.outcome,
1029                            ))
1030                        })
1031                    }
1032            }),
1033        )
1034        .await?;
1035
1036        Ok(WorkflowUpdateHandle::new(
1037            self.client.clone(),
1038            output.update_id,
1039            output.workflow_id,
1040            output.run_id.or_else(|| self.info().run_id.clone()),
1041            output.known_outcome,
1042        ))
1043    }
1044
1045    /// Get a handle to an existing update.
1046    ///
1047    /// The update definition determines the result type. The returned handle uses this workflow
1048    /// handle's workflow and run IDs and does not validate the update ID until
1049    /// [`get_result`](WorkflowUpdateHandle::get_result) is called.
1050    pub fn get_update_handle<U>(
1051        &self,
1052        update: U,
1053        update_id: impl Into<String>,
1054    ) -> WorkflowUpdateHandle<CT, U::Output>
1055    where
1056        U: UpdateDefinition<Workflow = W::Run>,
1057    {
1058        let _ = update;
1059        WorkflowUpdateHandle::new(
1060            self.client.clone(),
1061            update_id.into(),
1062            self.info.workflow_id.clone(),
1063            self.info.run_id.clone(),
1064            None,
1065        )
1066    }
1067
1068    /// Request cancellation of this workflow.
1069    pub async fn cancel(&self, opts: WorkflowCancelOptions) -> Result<(), WorkflowInteractionError>
1070    where
1071        CT: NamespacedClient,
1072    {
1073        interceptors::call_cancel_workflow(
1074            self.client.client_interceptors(),
1075            CancelWorkflowInput {
1076                workflow_id: self.info.workflow_id.clone(),
1077                run_id: self.info.run_id.clone().unwrap_or_default(),
1078                first_execution_run_id: self
1079                    .info
1080                    .first_execution_run_id
1081                    .clone()
1082                    .unwrap_or_default(),
1083                options: opts,
1084            },
1085            Next::new({
1086                let mut client = self.client.clone();
1087                move |input: CancelWorkflowInput| -> BoxFuture<
1088                    '_,
1089                    Result<(), WorkflowInteractionError>,
1090                > {
1091                    Box::pin(async move {
1092                        let mut request = RequestCancelWorkflowExecutionRequest {
1093                            namespace: client.namespace(),
1094                            workflow_execution: Some(ProtoWorkflowExecution {
1095                                workflow_id: input.workflow_id,
1096                                run_id: input.run_id,
1097                            }),
1098                            identity: client.identity(),
1099                            request_id: input
1100                                .options
1101                                .request_id
1102                                .clone()
1103                                .unwrap_or_else(|| Uuid::new_v4().to_string()),
1104                            first_execution_run_id: input.first_execution_run_id,
1105                            reason: input.options.reason.clone(),
1106                            links: vec![],
1107                        }
1108                        .into_request();
1109                        input.options.rpc_options.apply_to(&mut request);
1110                        WorkflowService::request_cancel_workflow_execution(&mut client, request)
1111                            .await
1112                            .map_err(WorkflowInteractionError::from_status)?;
1113                        Ok(())
1114                    })
1115                }
1116            }),
1117        )
1118        .await
1119    }
1120
1121    /// Terminate this workflow.
1122    pub async fn terminate(
1123        &self,
1124        opts: WorkflowTerminateOptions,
1125    ) -> Result<(), WorkflowInteractionError>
1126    where
1127        CT: NamespacedClient,
1128    {
1129        interceptors::call_terminate_workflow(
1130            self.client.client_interceptors(),
1131            TerminateWorkflowInput {
1132                workflow_id: self.info.workflow_id.clone(),
1133                run_id: self.info.run_id.clone().unwrap_or_default(),
1134                first_execution_run_id: self
1135                    .info
1136                    .first_execution_run_id
1137                    .clone()
1138                    .unwrap_or_default(),
1139                options: opts,
1140            },
1141            Next::new({
1142                let mut client = self.client.clone();
1143                move |input: TerminateWorkflowInput| -> BoxFuture<
1144                    '_,
1145                    Result<(), WorkflowInteractionError>,
1146                > {
1147                    Box::pin(async move {
1148                        let mut request = TerminateWorkflowExecutionRequest {
1149                            namespace: client.namespace(),
1150                            workflow_execution: Some(ProtoWorkflowExecution {
1151                                workflow_id: input.workflow_id,
1152                                run_id: input.run_id,
1153                            }),
1154                            reason: input.options.reason.clone(),
1155                            details: input.options.details.clone(),
1156                            identity: client.identity(),
1157                            first_execution_run_id: input.first_execution_run_id,
1158                            links: vec![],
1159                        }
1160                        .into_request();
1161                        input.options.rpc_options.apply_to(&mut request);
1162                        WorkflowService::terminate_workflow_execution(&mut client, request)
1163                            .await
1164                            .map_err(WorkflowInteractionError::from_status)?;
1165                        Ok(())
1166                    })
1167                }
1168            }),
1169        )
1170        .await
1171    }
1172
1173    /// Get workflow execution description/metadata.
1174    pub async fn describe(
1175        &self,
1176        opts: WorkflowDescribeOptions,
1177    ) -> Result<WorkflowExecutionDescription, WorkflowInteractionError>
1178    where
1179        CT: NamespacedClient,
1180    {
1181        let output = interceptors::call_describe_workflow(
1182            self.client.client_interceptors(),
1183            DescribeWorkflowInput {
1184                workflow_id: self.info.workflow_id.clone(),
1185                run_id: self.info.run_id.clone().unwrap_or_default(),
1186                options: opts,
1187            },
1188            Next::new({
1189                let mut client = self.client.clone();
1190                move |input: DescribeWorkflowInput| -> BoxFuture<
1191                        '_,
1192                        Result<DescribeWorkflowOutput, WorkflowInteractionError>,
1193                    > {
1194                        Box::pin(async move {
1195                            let mut request = DescribeWorkflowExecutionRequest {
1196                                namespace: client.namespace(),
1197                                execution: Some(ProtoWorkflowExecution {
1198                                    workflow_id: input.workflow_id,
1199                                    run_id: input.run_id,
1200                                }),
1201                            }
1202                            .into_request();
1203                            input.options.rpc_options.apply_to(&mut request);
1204                            let response =
1205                                WorkflowService::describe_workflow_execution(&mut client, request)
1206                                    .await
1207                                    .map_err(WorkflowInteractionError::from_status)?
1208                                    .into_inner();
1209                            Ok(DescribeWorkflowOutput::new(response))
1210                        })
1211                    }
1212            }),
1213        )
1214        .await?;
1215        WorkflowExecutionDescription::new(output.response, self.client.data_converter())
1216            .await
1217            .map_err(WorkflowInteractionError::from)
1218    }
1219    /// Fetch workflow execution history as a lazy stream.
1220    ///
1221    /// No request is sent until the returned stream is polled.
1222    pub fn fetch_history(&self, opts: WorkflowFetchHistoryOptions) -> WorkflowHistory
1223    where
1224        CT: NamespacedClient + 'static,
1225    {
1226        let run_id = self.info.run_id.clone().unwrap_or_default();
1227        self.fetch_history_for_run(&run_id, opts)
1228    }
1229
1230    fn fetch_history_for_run(
1231        &self,
1232        run_id: &str,
1233        opts: WorkflowFetchHistoryOptions,
1234    ) -> WorkflowHistory
1235    where
1236        CT: NamespacedClient + 'static,
1237    {
1238        let client = self.client.clone();
1239        let workflow_id = self.info.workflow_id.clone();
1240        let history_workflow_id = workflow_id.clone();
1241        let run_id = run_id.to_string();
1242
1243        let stream = stream::unfold(
1244            (Vec::new(), VecDeque::new(), false),
1245            move |(mut next_page_token, mut buffer, mut exhausted)| {
1246                let client = client.clone();
1247                let workflow_id = workflow_id.clone();
1248                let run_id = run_id.clone();
1249                let opts = opts.clone();
1250
1251                async move {
1252                    loop {
1253                        if let Some(event) = buffer.pop_front() {
1254                            return Some((Ok(event), (next_page_token, buffer, exhausted)));
1255                        }
1256
1257                        if exhausted {
1258                            return None;
1259                        }
1260
1261                        let output = interceptors::call_fetch_workflow_history_page(
1262                            client.client_interceptors(),
1263                            FetchWorkflowHistoryPageInput {
1264                                workflow_id: workflow_id.clone(),
1265                                run_id: run_id.clone(),
1266                                next_page_token: next_page_token.clone(),
1267                                options: opts.clone(),
1268                            },
1269                            Next::new({
1270                                let mut rpc_client = client.clone();
1271                                move |input: FetchWorkflowHistoryPageInput| -> BoxFuture<
1272                                    '_,
1273                                    Result<
1274                                        FetchWorkflowHistoryPageOutput,
1275                                        WorkflowInteractionError,
1276                                    >,
1277                                > {
1278                                    Box::pin(async move {
1279                                        let mut request = GetWorkflowExecutionHistoryRequest {
1280                                            namespace: rpc_client.namespace(),
1281                                            execution: Some(ProtoWorkflowExecution {
1282                                                workflow_id: input.workflow_id,
1283                                                run_id: input.run_id,
1284                                            }),
1285                                            next_page_token: input.next_page_token,
1286                                            skip_archival: input.options.skip_archival,
1287                                            wait_new_event: input.options.wait_new_event,
1288                                            history_event_filter_type:
1289                                                ProtoHistoryEventFilterType::from(
1290                                                    input.options.event_filter_type,
1291                                                )
1292                                                    as i32,
1293                                            ..Default::default()
1294                                        }
1295                                        .into_request();
1296                                        input.options.rpc_options.apply_to(&mut request);
1297                                        let response =
1298                                            WorkflowService::get_workflow_execution_history(
1299                                                &mut rpc_client,
1300                                                request,
1301                                            )
1302                                            .await
1303                                            .map_err(WorkflowInteractionError::from_status)?
1304                                            .into_inner();
1305                                        Ok(FetchWorkflowHistoryPageOutput::new(
1306                                            response
1307                                                .history
1308                                                .map(|history| history.events)
1309                                                .unwrap_or_default(),
1310                                            response.next_page_token,
1311                                        ))
1312                                    })
1313                                }
1314                            }),
1315                        )
1316                        .await;
1317
1318                        match output {
1319                            Ok(output) => {
1320                                exhausted = output.next_page_token.is_empty();
1321                                next_page_token = output.next_page_token;
1322                                buffer = output.events.into();
1323                            }
1324                            Err(error) => {
1325                                return Some((Err(error), (next_page_token, buffer, true)));
1326                            }
1327                        }
1328                    }
1329                }
1330            },
1331        );
1332
1333        WorkflowHistory {
1334            inner: Box::pin(stream),
1335            workflow_id: Some(history_workflow_id),
1336        }
1337    }
1338}
1339
1340/// Handle to a workflow update that has been started but may not be complete.
1341///
1342/// Use [`get_result`](Self::get_result) to wait for the update to complete and retrieve its result.
1343pub struct WorkflowUpdateHandle<CT, T> {
1344    client: CT,
1345    update_id: String,
1346    workflow_id: String,
1347    run_id: Option<String>,
1348    /// If the update was started with `Completed` wait stage, the outcome is already available.
1349    known_outcome: Option<update::v1::Outcome>,
1350    _output: PhantomData<T>,
1351}
1352
1353impl<CT, T> WorkflowUpdateHandle<CT, T> {
1354    pub(crate) fn new(
1355        client: CT,
1356        update_id: String,
1357        workflow_id: String,
1358        run_id: Option<String>,
1359        known_outcome: Option<update::v1::Outcome>,
1360    ) -> Self {
1361        Self {
1362            client,
1363            update_id,
1364            workflow_id,
1365            run_id,
1366            known_outcome,
1367            _output: PhantomData,
1368        }
1369    }
1370
1371    /// Get the update ID.
1372    pub fn id(&self) -> &str {
1373        &self.update_id
1374    }
1375
1376    /// Get the workflow ID.
1377    pub fn workflow_id(&self) -> &str {
1378        &self.workflow_id
1379    }
1380
1381    /// Get the workflow run ID, if available.
1382    pub fn workflow_run_id(&self) -> Option<&str> {
1383        self.run_id.as_deref()
1384    }
1385}
1386
1387impl<CT, T: 'static> WorkflowUpdateHandle<CT, T>
1388where
1389    CT: WorkflowService + NamespacedClient + Clone,
1390{
1391    /// Wait for the update to complete and return the result using the provided RPC controls.
1392    pub async fn get_result(&self, rpc_options: RpcOptions) -> Result<T, WorkflowUpdateError>
1393    where
1394        T: temporalio_common::data_converters::TemporalDeserializable,
1395    {
1396        let output = interceptors::call_poll_workflow_update(
1397            self.client.client_interceptors(),
1398            PollWorkflowUpdateInput {
1399                update_id: self.update_id.clone(),
1400                workflow_id: self.workflow_id.clone(),
1401                run_id: self.run_id.clone().unwrap_or_default(),
1402                rpc_options,
1403            },
1404            Next::new({
1405                let mut client = self.client.clone();
1406                let known_outcome = self.known_outcome.clone();
1407                move |input: PollWorkflowUpdateInput| -> BoxFuture<
1408                    '_,
1409                    Result<PollWorkflowUpdateOutput, WorkflowUpdateError>,
1410                > {
1411                    Box::pin(async move {
1412                        if let Some(outcome) = known_outcome {
1413                            return Ok(PollWorkflowUpdateOutput::new(outcome));
1414                        }
1415                        // The server's internal long-poll timeout (~60s) may expire before the update
1416                        // completes, returning a response with outcome: None. Keep polling until we
1417                        // get an actual outcome.
1418                        loop {
1419                            let mut request = PollWorkflowExecutionUpdateRequest {
1420                                namespace: client.namespace(),
1421                                update_ref: Some(update::v1::UpdateRef {
1422                                    workflow_execution: Some(ProtoWorkflowExecution {
1423                                        workflow_id: input.workflow_id.clone(),
1424                                        run_id: input.run_id.clone(),
1425                                    }),
1426                                    update_id: input.update_id.clone(),
1427                                }),
1428                                identity: client.identity(),
1429                                wait_policy: Some(WaitPolicy {
1430                                    lifecycle_stage:
1431                                        UpdateWorkflowExecutionLifecycleStage::Completed.into(),
1432                                }),
1433                            }
1434                            .into_request();
1435                            input.rpc_options.apply_to(&mut request);
1436                            let response = WorkflowService::poll_workflow_execution_update(
1437                                &mut client,
1438                                request,
1439                            )
1440                            .await
1441                            .map_err(WorkflowUpdateError::from_status)?
1442                            .into_inner();
1443                            if let Some(outcome) = response.outcome {
1444                                return Ok(PollWorkflowUpdateOutput::new(outcome));
1445                            }
1446                        }
1447                    })
1448                }
1449            }),
1450        )
1451        .await?;
1452        let outcome = output.outcome;
1453
1454        match outcome.value {
1455            Some(update::v1::outcome::Value::Success(success)) => self
1456                .client
1457                .data_converter()
1458                .from_payloads(
1459                    &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1460                    success.payloads,
1461                )
1462                .await
1463                .map_err(WorkflowUpdateError::from),
1464            Some(update::v1::outcome::Value::Failure(failure)) => {
1465                Err(WorkflowUpdateError::Failed(Box::new(failure)))
1466            }
1467            None => Err(WorkflowUpdateError::Other(
1468                "Update returned no outcome value".into(),
1469            )),
1470        }
1471    }
1472}
1473
1474#[cfg(test)]
1475mod tests {
1476    use super::*;
1477    use crate::{ClientInterceptor, test_helpers::XorCodec};
1478    use futures_util::{FutureExt, StreamExt};
1479    use std::{
1480        collections::{HashMap, VecDeque},
1481        sync::{
1482            Arc, Mutex,
1483            atomic::{AtomicUsize, Ordering},
1484        },
1485    };
1486    use temporalio_common::{
1487        data_converters::DefaultFailureConverter,
1488        protos::temporal::api::{
1489            common::v1::{Memo, SearchAttributes, WorkflowExecution},
1490            enums::v1::{
1491                UpdateWorkflowExecutionLifecycleStage as UpdateStage,
1492                WorkflowExecutionStatus as ProtoWorkflowExecutionStatus,
1493            },
1494            history::v1::WorkflowExecutionStartedEventAttributes,
1495            sdk::v1::UserMetadata,
1496            update::v1::UpdateRef,
1497            workflow::v1::WorkflowExecutionConfig,
1498            workflowservice::v1::{
1499                GetWorkflowExecutionHistoryResponse, UpdateWorkflowExecutionResponse,
1500            },
1501        },
1502    };
1503    use tonic::{Request, Response, Status};
1504
1505    #[tokio::test]
1506    async fn workflow_history_workflow_id_roundtrips() {
1507        let event = HistoryEvent {
1508            event_id: 1,
1509            attributes: Some(Attributes::WorkflowExecutionStartedEventAttributes(
1510                WorkflowExecutionStartedEventAttributes {
1511                    workflow_id: "workflow-id".to_owned(),
1512                    original_execution_run_id: "run-id".to_owned(),
1513                    ..Default::default()
1514                },
1515            )),
1516            ..Default::default()
1517        };
1518        let history = WorkflowHistory {
1519            inner: Box::pin(stream::iter(std::iter::once(Ok(event)))),
1520            workflow_id: None,
1521        };
1522
1523        let bytes = history.to_json().await.unwrap();
1524
1525        let decoded = WorkflowHistory::from_json(&bytes).unwrap();
1526        assert_eq!(decoded.workflow_id(), Some("workflow-id"));
1527    }
1528
1529    #[derive(Clone)]
1530    struct MockHistoryClient {
1531        responses: Arc<Mutex<VecDeque<Result<GetWorkflowExecutionHistoryResponse, tonic::Status>>>>,
1532        calls: Arc<AtomicUsize>,
1533        interceptors: Vec<Arc<dyn ClientInterceptor>>,
1534    }
1535
1536    impl NamespacedClient for MockHistoryClient {
1537        fn namespace(&self) -> String {
1538            "test-namespace".to_owned()
1539        }
1540
1541        fn identity(&self) -> String {
1542            "test-identity".to_owned()
1543        }
1544
1545        fn client_interceptors(&self) -> &[Arc<dyn ClientInterceptor>] {
1546            &self.interceptors
1547        }
1548    }
1549
1550    impl WorkflowService for MockHistoryClient {
1551        fn get_workflow_execution_history(
1552            &mut self,
1553            _request: Request<GetWorkflowExecutionHistoryRequest>,
1554        ) -> BoxFuture<'_, Result<Response<GetWorkflowExecutionHistoryResponse>, tonic::Status>>
1555        {
1556            self.calls.fetch_add(1, Ordering::SeqCst);
1557            let response = self.responses.lock().unwrap().pop_front().unwrap();
1558            async move { response.map(Response::new) }.boxed()
1559        }
1560    }
1561
1562    struct CountingHistoryInterceptor(Arc<AtomicUsize>);
1563
1564    impl ClientInterceptor for CountingHistoryInterceptor {
1565        fn fetch_workflow_history_page<'a>(
1566            &'a self,
1567            input: FetchWorkflowHistoryPageInput,
1568            next: Next<
1569                'a,
1570                FetchWorkflowHistoryPageInput,
1571                BoxFuture<'a, Result<FetchWorkflowHistoryPageOutput, WorkflowInteractionError>>,
1572            >,
1573        ) -> BoxFuture<'a, Result<FetchWorkflowHistoryPageOutput, WorkflowInteractionError>>
1574        {
1575            self.0.fetch_add(1, Ordering::SeqCst);
1576            next.run(input)
1577        }
1578    }
1579
1580    fn history_response(
1581        event_ids: impl IntoIterator<Item = i64>,
1582        next_page_token: &[u8],
1583    ) -> GetWorkflowExecutionHistoryResponse {
1584        GetWorkflowExecutionHistoryResponse {
1585            history: Some(History {
1586                events: event_ids
1587                    .into_iter()
1588                    .map(|event_id| HistoryEvent {
1589                        event_id,
1590                        ..Default::default()
1591                    })
1592                    .collect(),
1593            }),
1594            next_page_token: next_page_token.to_vec(),
1595            ..Default::default()
1596        }
1597    }
1598
1599    fn history_handle(
1600        responses: impl IntoIterator<Item = Result<GetWorkflowExecutionHistoryResponse, tonic::Status>>,
1601        calls: Arc<AtomicUsize>,
1602        interceptors: Vec<Arc<dyn ClientInterceptor>>,
1603    ) -> WorkflowHandle<MockHistoryClient, UntypedWorkflow> {
1604        WorkflowHandle::new(
1605            MockHistoryClient {
1606                responses: Arc::new(Mutex::new(responses.into_iter().collect())),
1607                calls,
1608                interceptors,
1609            },
1610            WorkflowExecutionInfo {
1611                namespace: "test-namespace".to_owned(),
1612                workflow_id: "workflow-id".to_owned(),
1613                run_id: Some("run-id".to_owned()),
1614                first_execution_run_id: None,
1615            },
1616        )
1617    }
1618
1619    #[tokio::test]
1620    async fn workflow_history_fetches_pages_lazily() {
1621        let calls = Arc::new(AtomicUsize::new(0));
1622        let interceptor_calls = Arc::new(AtomicUsize::new(0));
1623        let handle = history_handle(
1624            [
1625                Ok(history_response([], b"second-page")),
1626                Ok(history_response([1, 2], b"third-page")),
1627                Ok(history_response([3], b"")),
1628            ],
1629            calls.clone(),
1630            vec![Arc::new(CountingHistoryInterceptor(
1631                interceptor_calls.clone(),
1632            ))],
1633        );
1634
1635        let mut history = handle.fetch_history(WorkflowFetchHistoryOptions::default());
1636        assert_eq!(calls.load(Ordering::SeqCst), 0);
1637
1638        assert_eq!(history.next().await.unwrap().unwrap().event_id, 1);
1639        assert_eq!(calls.load(Ordering::SeqCst), 2);
1640        assert_eq!(history.next().await.unwrap().unwrap().event_id, 2);
1641        assert_eq!(calls.load(Ordering::SeqCst), 2);
1642        assert_eq!(history.next().await.unwrap().unwrap().event_id, 3);
1643        assert_eq!(calls.load(Ordering::SeqCst), 3);
1644        assert!(history.next().await.is_none());
1645        assert_eq!(interceptor_calls.load(Ordering::SeqCst), 3);
1646    }
1647
1648    #[tokio::test]
1649    async fn workflow_history_yields_page_error_then_ends() {
1650        let calls = Arc::new(AtomicUsize::new(0));
1651        let handle = history_handle(
1652            [
1653                Ok(history_response([1], b"second-page")),
1654                Err(tonic::Status::unavailable("history unavailable")),
1655            ],
1656            calls.clone(),
1657            Vec::new(),
1658        );
1659        let mut history = handle.fetch_history(WorkflowFetchHistoryOptions::default());
1660
1661        assert_eq!(history.next().await.unwrap().unwrap().event_id, 1);
1662        assert!(matches!(
1663            history.next().await.unwrap(),
1664            Err(WorkflowInteractionError::Rpc(status)) if status.code() == tonic::Code::Unavailable
1665        ));
1666        assert!(history.next().await.is_none());
1667        assert_eq!(calls.load(Ordering::SeqCst), 2);
1668    }
1669
1670    #[tokio::test]
1671    async fn workflow_result_details_support_typed_decoding() {
1672        let converter = DataConverter::new(
1673            PayloadConverter::default(),
1674            DefaultFailureConverter::default(),
1675            XorCodec,
1676        );
1677        let payloads = converter
1678            .to_payloads(
1679                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1680                &"workflow-result-details".to_owned(),
1681            )
1682            .await
1683            .unwrap();
1684        let details = WorkflowResultDetails::new(payloads.clone(), &converter)
1685            .await
1686            .unwrap();
1687
1688        assert_ne!(details.raw(), payloads);
1689        let decoded_payloads = details.raw().to_vec();
1690        assert_eq!(
1691            details.deserialize::<String>().unwrap(),
1692            "workflow-result-details"
1693        );
1694        assert_eq!(details.into_raw().payloads, decoded_payloads);
1695    }
1696
1697    #[tokio::test]
1698    async fn workflow_result_detail_conversion_errors_are_reported() {
1699        let details =
1700            WorkflowResultDetails::new(vec![Payload::default()], &DataConverter::default())
1701                .await
1702                .unwrap();
1703
1704        assert_eq!(details.raw(), &[Payload::default()]);
1705        assert!(details.deserialize::<String>().is_err());
1706    }
1707
1708    #[tokio::test]
1709    async fn workflow_description_memo_uses_saved_converter() {
1710        let converter = DataConverter::new(
1711            PayloadConverter::default(),
1712            DefaultFailureConverter::default(),
1713            XorCodec,
1714        );
1715        let encoded = converter
1716            .to_payload(
1717                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1718                &"memo-value".to_owned(),
1719            )
1720            .await
1721            .unwrap();
1722        let description = WorkflowExecutionDescription::new(
1723            DescribeWorkflowExecutionResponse {
1724                workflow_execution_info: Some(workflow::WorkflowExecutionInfo {
1725                    memo: Some(Memo {
1726                        fields: HashMap::from([("memo-key".to_owned(), encoded)]),
1727                    }),
1728                    ..Default::default()
1729                }),
1730                ..Default::default()
1731            },
1732            &converter,
1733        )
1734        .await
1735        .unwrap();
1736        let memo = description.memo();
1737
1738        assert_eq!(
1739            memo.get::<String>("memo-key").unwrap(),
1740            Some("memo-value".to_owned())
1741        );
1742    }
1743
1744    #[tokio::test]
1745    async fn workflow_description_accessors_expose_decoded_fields() {
1746        let converter = DataConverter::default();
1747        let memo_payload = converter
1748            .to_payload(
1749                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1750                &"memo-value",
1751            )
1752            .await
1753            .unwrap();
1754        let search_attr_payload = converter
1755            .to_payload(
1756                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1757                &"search-value",
1758            )
1759            .await
1760            .unwrap();
1761        let summary_payload = converter
1762            .to_payload(
1763                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1764                &"workflow summary",
1765            )
1766            .await
1767            .unwrap();
1768        let details_payload = converter
1769            .to_payload(
1770                &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1771                &"workflow details",
1772            )
1773            .await
1774            .unwrap();
1775        let description = WorkflowExecutionDescription::new(
1776            DescribeWorkflowExecutionResponse {
1777                workflow_execution_info: Some(workflow::WorkflowExecutionInfo {
1778                    execution: Some(ProtoWorkflowExecution {
1779                        workflow_id: "wf-id".to_string(),
1780                        run_id: "run-id".to_string(),
1781                    }),
1782                    r#type: Some(
1783                        temporalio_common::protos::temporal::api::common::v1::WorkflowType {
1784                            name: "wf-type".to_string(),
1785                        },
1786                    ),
1787                    status: ProtoWorkflowExecutionStatus::Completed as i32,
1788                    task_queue: "task-queue".to_string(),
1789                    history_length: 42,
1790                    memo: Some(Memo {
1791                        fields: HashMap::from([("memo-key".to_string(), memo_payload.clone())]),
1792                    }),
1793                    parent_execution: Some(ProtoWorkflowExecution {
1794                        workflow_id: "parent-id".to_string(),
1795                        run_id: "parent-run-id".to_string(),
1796                    }),
1797                    search_attributes: Some(SearchAttributes {
1798                        indexed_fields: HashMap::from([(
1799                            "CustomKeywordField".to_string(),
1800                            search_attr_payload.clone(),
1801                        )]),
1802                    }),
1803                    ..Default::default()
1804                }),
1805                execution_config: Some(WorkflowExecutionConfig {
1806                    user_metadata: Some(UserMetadata {
1807                        summary: Some(summary_payload),
1808                        details: Some(details_payload),
1809                    }),
1810                    ..Default::default()
1811                }),
1812                ..Default::default()
1813            },
1814            &converter,
1815        )
1816        .await
1817        .unwrap();
1818
1819        assert_eq!(description.id(), "wf-id");
1820        assert_eq!(description.run_id(), "run-id");
1821        assert_eq!(description.workflow_type(), "wf-type");
1822        assert_eq!(description.status(), WorkflowExecutionStatus::Completed);
1823        let mut unknown_status_description = description.clone();
1824        unknown_status_description
1825            .raw_description
1826            .workflow_execution_info
1827            .as_mut()
1828            .unwrap()
1829            .status = 123_456;
1830        assert_eq!(
1831            unknown_status_description.status(),
1832            WorkflowExecutionStatus::Unknown
1833        );
1834        assert_eq!(description.task_queue(), "task-queue");
1835        assert_eq!(description.history_length(), 42);
1836        assert_eq!(description.parent_id(), Some("parent-id"));
1837        assert_eq!(description.parent_run_id(), Some("parent-run-id"));
1838        let memo = description.memo();
1839        assert_eq!(memo.raw_value("memo-key"), Some(&memo_payload));
1840        assert_eq!(
1841            memo.get::<String>("memo-key").unwrap(),
1842            Some("memo-value".to_owned())
1843        );
1844        let search_attributes = description.search_attributes();
1845        assert_eq!(
1846            search_attributes.raw_payload("CustomKeywordField"),
1847            Some(&search_attr_payload)
1848        );
1849        assert_eq!(description.static_summary(), Some("workflow summary"));
1850        assert_eq!(description.static_details(), Some("workflow details"));
1851    }
1852
1853    #[tokio::test]
1854    async fn workflow_description_rejects_negative_history_length() {
1855        let err = WorkflowExecutionDescription::new(
1856            DescribeWorkflowExecutionResponse {
1857                workflow_execution_info: Some(workflow::WorkflowExecutionInfo {
1858                    history_length: -1,
1859                    ..Default::default()
1860                }),
1861                ..Default::default()
1862            },
1863            &DataConverter::default(),
1864        )
1865        .await
1866        .unwrap_err();
1867
1868        assert_eq!(
1869            err.to_string(),
1870            "Encoding error: workflow history_length must be non-negative, got -1"
1871        );
1872    }
1873
1874    #[derive(Default)]
1875    struct MockUpdateState {
1876        responses: VecDeque<Result<UpdateWorkflowExecutionResponse, Status>>,
1877        requests: Vec<UpdateWorkflowExecutionRequest>,
1878    }
1879
1880    #[derive(Clone)]
1881    struct MockUpdateClient(Arc<Mutex<MockUpdateState>>);
1882
1883    impl NamespacedClient for MockUpdateClient {
1884        fn namespace(&self) -> String {
1885            "ns".into()
1886        }
1887        fn identity(&self) -> String {
1888            "identity".into()
1889        }
1890    }
1891
1892    impl WorkflowService for MockUpdateClient {
1893        fn update_workflow_execution(
1894            &mut self,
1895            request: Request<UpdateWorkflowExecutionRequest>,
1896        ) -> BoxFuture<'_, Result<Response<UpdateWorkflowExecutionResponse>, Status>> {
1897            let mut state = self.0.lock().unwrap();
1898            state.requests.push(request.into_inner());
1899            let response = state.responses.pop_front().expect("unexpected submission");
1900            Box::pin(async { response.map(Response::new) })
1901        }
1902    }
1903
1904    #[rstest::rstest]
1905    #[case::retries_until_accepted(
1906    vec![UpdateStage::Unspecified, UpdateStage::Admitted, UpdateStage::Admitted, UpdateStage::Accepted],
1907    None
1908)]
1909    #[case::already_accepted(vec![UpdateStage::Accepted], None)]
1910    #[case::already_completed(vec![UpdateStage::Completed], None)]
1911    #[case::rpc_error(vec![UpdateStage::Admitted], Some(Status::unavailable("transport failure")))]
1912    #[tokio::test]
1913    async fn update_waits_for_acceptance(
1914        #[case] stages: Vec<UpdateStage>,
1915        #[case] error: Option<Status>,
1916    ) {
1917        use std::time::Duration;
1918
1919        let mut responses: VecDeque<_> = stages
1920            .iter()
1921            .map(|stage| {
1922                Ok(UpdateWorkflowExecutionResponse {
1923                    stage: *stage as i32,
1924                    update_ref: Some(UpdateRef {
1925                        workflow_execution: Some(WorkflowExecution {
1926                            workflow_id: "wf".into(),
1927                            run_id: "run".into(),
1928                        }),
1929                        ..Default::default()
1930                    }),
1931                    ..Default::default()
1932                })
1933            })
1934            .collect();
1935        // Only successful responses below Accepted should trigger another submission.
1936        if let Some(error) = error.clone() {
1937            responses.push_back(Err(error));
1938        }
1939        let expected_requests = responses.len();
1940        let state = Arc::new(Mutex::new(MockUpdateState {
1941            responses,
1942            ..Default::default()
1943        }));
1944        let handle = WorkflowHandle::<_, UntypedWorkflow>::new(
1945            MockUpdateClient(state.clone()),
1946            WorkflowExecutionInfo::builder()
1947                .namespace("ns")
1948                .workflow_id("wf")
1949                .build(),
1950        );
1951        let result = handle
1952            .start_update(
1953                UntypedUpdate::new("handler"),
1954                RawValue::new(vec![Payload {
1955                    data: vec![1, 2, 3],
1956                    ..Default::default()
1957                }]),
1958                WorkflowStartUpdateOptions::builder()
1959                    .rpc_options(
1960                        RpcOptions::builder()
1961                            .timeout(Duration::from_secs(30))
1962                            .build(),
1963                    )
1964                    .build(),
1965            )
1966            .await;
1967        let state = state.lock().unwrap();
1968        assert_eq!(state.requests.len(), expected_requests);
1969        let first = &state.requests[0];
1970        // A retry must remain the same logical update
1971        assert!(state.requests.iter().all(|request| request == first));
1972        assert!(first.workflow_execution.as_ref().unwrap().run_id.is_empty());
1973        let update_id = &first
1974            .request
1975            .as_ref()
1976            .unwrap()
1977            .meta
1978            .as_ref()
1979            .unwrap()
1980            .update_id;
1981        assert!(!update_id.is_empty());
1982        match (error, result) {
1983            (None, Ok(handle)) => {
1984                assert_eq!(handle.id(), update_id);
1985                assert_eq!(handle.workflow_run_id(), Some("run"));
1986            }
1987            (Some(expected), Err(WorkflowUpdateError::Rpc(actual))) => {
1988                assert_eq!(actual.code(), expected.code());
1989                assert_eq!(actual.message(), expected.message());
1990            }
1991            _ => panic!("unexpected update result"),
1992        }
1993    }
1994}