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#[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 pub fn deserialize<T: temporalio_common::data_converters::TemporalDeserializable + 'static>(
123 &self,
124 ) -> Result<T, PayloadConversionError> {
125 self.payloads.deserialize()
126 }
127
128 pub fn raw(&self) -> &[Payload] {
130 self.payloads.raw()
131 }
132
133 pub fn into_raw(self) -> RawValue {
135 self.payloads.into_raw()
136 }
137}
138
139#[derive(Debug)]
141#[allow(clippy::large_enum_variant)]
142pub enum WorkflowExecutionResult<T> {
143 Succeeded(T),
145 Failed(IncomingError),
147 Cancelled {
149 details: WorkflowResultDetails,
151 },
152 Terminated {
154 details: WorkflowResultDetails,
156 },
157 TimedOut,
159 ContinuedAsNew,
161}
162
163#[derive(Debug, Clone)]
167pub struct WorkflowExecutionDescription {
168 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 pub fn id(&self) -> &str {
217 self.execution().workflow_id.as_str()
218 }
219
220 pub fn run_id(&self) -> &str {
222 self.execution().run_id.as_str()
223 }
224
225 pub fn workflow_type(&self) -> &str {
227 self.workflow_type_info().name.as_str()
228 }
229
230 pub fn status(&self) -> WorkflowExecutionStatus {
232 WorkflowExecutionStatus::from_raw(self.workflow_info().status)
233 }
234
235 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 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 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 pub fn task_queue(&self) -> &str {
261 self.workflow_info().task_queue.as_str()
262 }
263
264 pub fn history_length(&self) -> usize {
266 self.history_length
267 }
268
269 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 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 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 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 pub fn static_summary(&self) -> Option<&str> {
305 self.static_summary.as_deref()
306 }
307
308 pub fn static_details(&self) -> Option<&str> {
310 self.static_details.as_deref()
311 }
312
313 pub fn raw(&self) -> &DescribeWorkflowExecutionResponse {
315 &self.raw_description
316 }
317
318 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#[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#[derive(Debug, thiserror::Error)]
389#[non_exhaustive]
390pub enum WorkflowHistoryError {
391 #[error("failed to fetch workflow history: {0}")]
393 Fetch(#[from] WorkflowInteractionError),
394 #[error("failed to convert workflow history JSON: {0}")]
396 Json(#[from] serde_json::Error),
397}
398
399impl WorkflowHistory {
400 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 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 pub fn workflow_id(&self) -> Option<&str> {
415 self.workflow_id.as_deref()
416 }
417
418 pub async fn into_events(self) -> Result<Vec<HistoryEvent>, WorkflowInteractionError> {
420 self.inner.try_collect().await
421 }
422}
423
424#[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 pub fn run_id(&self) -> Option<&str> {
437 self.info.run_id.as_deref()
438 }
439}
440
441#[derive(Debug, Clone, bon::Builder)]
443#[builder(on(String, into), state_mod(vis = "pub"))]
444#[non_exhaustive]
445pub struct WorkflowExecutionInfo {
446 pub namespace: String,
448 pub workflow_id: String,
450 pub run_id: Option<String>,
452 pub first_execution_run_id: Option<String>,
456}
457
458impl WorkflowExecutionInfo {
459 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
468pub type UntypedWorkflowHandle<CT> = WorkflowHandle<CT, UntypedWorkflow>;
471
472pub struct UntypedSignal<W> {
476 name: String,
477 _wf: PhantomData<W>,
478}
479
480impl<W> UntypedSignal<W> {
481 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
499pub struct UntypedQuery<W> {
503 name: String,
504 _wf: PhantomData<W>,
505}
506
507impl<W> UntypedQuery<W> {
508 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
527pub struct UntypedUpdate<W> {
531 name: String,
532 _wf: PhantomData<W>,
533}
534
535impl<W> UntypedUpdate<W> {
536 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#[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 pub fn new(client: CT, info: WorkflowExecutionInfo) -> Self {
601 Self {
602 client,
603 info,
604 _wf_type: PhantomData::<W>,
605 }
606 }
607
608 pub fn info(&self) -> &WorkflowExecutionInfo {
610 &self.info
611 }
612
613 pub fn client(&self) -> &CT {
615 &self.client
616 }
617
618 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 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 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 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 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 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 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 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 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 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 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
1340pub struct WorkflowUpdateHandle<CT, T> {
1344 client: CT,
1345 update_id: String,
1346 workflow_id: String,
1347 run_id: Option<String>,
1348 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 pub fn id(&self) -> &str {
1373 &self.update_id
1374 }
1375
1376 pub fn workflow_id(&self) -> &str {
1378 &self.workflow_id
1379 }
1380
1381 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 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 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 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 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}