1use super::{
13 GenericPayloadConverter, PayloadConversionError, PayloadConverter, SerializationContext,
14 SerializationContextData,
15};
16use crate::{
17 error::{
18 ActivityExecutionError, ActivityFailureError, ApplicationFailure, CancelledError,
19 ChildWorkflowExecutionError, ChildWorkflowFailureError, ChildWorkflowStartError,
20 IncomingError, IncomingNexusHandlerError, IncomingNexusOperationExecutionError,
21 OutgoingActivityError, OutgoingError, OutgoingWorkflowError, ResetWorkflowError,
22 ServerError, TerminatedError, TimeoutError, WorkflowSignalError,
23 WorkflowSignalFailureError,
24 },
25 protos::temporal::api::{
26 enums::v1::ApplicationErrorCategory as ProtoApplicationErrorCategory,
27 failure::v1::{
28 ActivityFailureInfo, ApplicationFailureInfo, CanceledFailureInfo,
29 ChildWorkflowExecutionFailureInfo, Failure, failure::FailureInfo,
30 },
31 },
32};
33
34pub trait FailureConverter {
36 fn to_failure(
38 &self,
39 error: OutgoingError,
40 payload_converter: &PayloadConverter,
41 context: &SerializationContextData,
42 ) -> Failure;
43
44 fn to_error(
46 &self,
47 failure: Failure,
48 payload_converter: &PayloadConverter,
49 context: &SerializationContextData,
50 ) -> Result<IncomingError, PayloadConversionError>;
51}
52
53pub struct DefaultFailureConverter {
55 encode_common_attributes: bool,
56}
57
58impl DefaultFailureConverter {
59 pub const fn new(encode_common_attributes: bool) -> Self {
62 Self {
63 encode_common_attributes,
64 }
65 }
66}
67
68impl Default for DefaultFailureConverter {
69 fn default() -> Self {
70 Self::new(false)
71 }
72}
73
74#[derive(serde::Deserialize, serde::Serialize)]
76#[non_exhaustive]
77pub struct CommonAttributes {
78 pub message: String,
80 pub stack_trace: String,
82}
83
84pub trait FailureDecodeHint {
86 type Output;
88
89 fn adapt(self, normalized: IncomingError) -> Self::Output;
91}
92
93#[derive(Debug, Clone, Copy)]
95pub struct NoopDecodeHint;
96
97impl FailureDecodeHint for NoopDecodeHint {
98 type Output = IncomingError;
99
100 fn adapt(self, normalized: IncomingError) -> Self::Output {
101 normalized
102 }
103}
104
105#[derive(Debug, Clone, Copy)]
107#[non_exhaustive]
108pub struct ActivityExecutionDecodeHint {
109 pub cancelled: bool,
111}
112
113impl ActivityExecutionDecodeHint {
114 pub fn new(cancelled: bool) -> Self {
116 Self { cancelled }
117 }
118}
119
120impl FailureDecodeHint for ActivityExecutionDecodeHint {
121 type Output = ActivityExecutionError;
122
123 fn adapt(self, normalized: IncomingError) -> Self::Output {
124 match normalized {
125 IncomingError::Activity(activity) => {
126 if self.cancelled && matches!(activity.cause(), Some(IncomingError::Cancelled(_))) {
127 let (_, cause) = activity.into_parts();
130 let Some(IncomingError::Cancelled(cancelled)) = cause else {
131 unreachable!("checked above");
132 };
133 ActivityExecutionError::Cancelled(cancelled)
134 } else {
135 ActivityExecutionError::Failed(activity)
136 }
137 }
138 other => match other {
139 IncomingError::Cancelled(cancelled) if self.cancelled => {
140 ActivityExecutionError::Cancelled(cancelled)
141 }
142 other => {
143 let activity = ActivityFailureError::new(
144 other.into_failure(),
145 ActivityFailureInfo::default(),
146 None,
147 );
148 ActivityExecutionError::Failed(activity)
149 }
150 },
151 }
152 }
153}
154
155#[derive(Debug, Clone, Copy, Default)]
157#[non_exhaustive]
158pub struct ChildWorkflowStartDecodeHint;
159
160impl FailureDecodeHint for ChildWorkflowStartDecodeHint {
161 type Output = ChildWorkflowStartError;
162
163 fn adapt(self, normalized: IncomingError) -> Self::Output {
164 match normalized {
165 IncomingError::Cancelled(cancelled) => {
166 ChildWorkflowStartError::Cancelled(Box::new(cancelled))
167 }
168 other => {
169 let payload_converter = PayloadConverter::default();
170 ChildWorkflowStartError::Cancelled(Box::new(CancelledError::new(
171 other.into_failure(),
172 CanceledFailureInfo::default(),
173 None,
174 &payload_converter,
175 &SerializationContextData::None,
176 )))
177 }
178 }
179 }
180}
181
182#[derive(Debug, Clone, Copy, Default)]
184#[non_exhaustive]
185pub struct ChildWorkflowExecutionDecodeHint;
186
187impl FailureDecodeHint for ChildWorkflowExecutionDecodeHint {
188 type Output = ChildWorkflowExecutionError;
189
190 fn adapt(self, normalized: IncomingError) -> Self::Output {
191 match normalized {
192 IncomingError::ChildWorkflowExecution(child) => {
193 ChildWorkflowExecutionError::Failed(Box::new(child))
194 }
195 other => ChildWorkflowExecutionError::Failed(Box::new(ChildWorkflowFailureError::new(
196 other.into_failure(),
197 ChildWorkflowExecutionFailureInfo::default(),
198 None,
199 ))),
200 }
201 }
202}
203
204#[derive(Debug, Clone, Copy, Default)]
206#[non_exhaustive]
207pub struct WorkflowSignalDecodeHint;
208
209impl FailureDecodeHint for WorkflowSignalDecodeHint {
210 type Output = WorkflowSignalError;
211
212 fn adapt(self, normalized: IncomingError) -> Self::Output {
213 let failure = normalized.failure().clone();
214 WorkflowSignalError::Failed(Box::new(WorkflowSignalFailureError::new(
215 failure, normalized,
216 )))
217 }
218}
219
220impl FailureConverter for DefaultFailureConverter {
221 fn to_failure(
222 &self,
223 error: OutgoingError,
224 payload_converter: &PayloadConverter,
225 context: &SerializationContextData,
226 ) -> Failure {
227 let original_error = error.to_string();
228 let encoded = match error {
229 OutgoingError::Activity(activity) => {
230 encode_outgoing_activity_error(activity, payload_converter, context)
231 }
232 OutgoingError::Workflow(OutgoingWorkflowError::Application(app)) => {
233 app.encode_failure(payload_converter, context)
234 }
235 OutgoingError::Workflow(OutgoingWorkflowError::PayloadConversion(err)) => {
236 Ok(encode_generic_application_failure(&err))
237 }
238 OutgoingError::Workflow(OutgoingWorkflowError::ActivityExecution(activity)) => {
239 activity.encode_failure(payload_converter, context)
240 }
241 OutgoingError::Workflow(OutgoingWorkflowError::ChildWorkflowExecution(child)) => {
242 child.encode_failure(payload_converter, context)
243 }
244 OutgoingError::Workflow(OutgoingWorkflowError::ChildWorkflowStart(child)) => {
245 child.encode_failure(payload_converter, context)
246 }
247 OutgoingError::Workflow(OutgoingWorkflowError::WorkflowSignal(signal)) => {
248 signal.encode_failure(payload_converter, context)
249 }
250 };
251 let mut failure = encoded.unwrap_or_else(|converter_error| {
252 Failure::application_failure(
253 failed_error_conversion_message(&original_error, &converter_error),
254 false,
255 )
256 });
257 if self.encode_common_attributes
258 && encode_common_attributes(&mut failure, payload_converter, context).is_err()
259 {
260 failure = Failure::application_failure(
261 "Failed encoding failure attributes".to_owned(),
262 false,
263 );
264 }
265 failure
266 }
267
268 fn to_error(
269 &self,
270 failure: Failure,
271 payload_converter: &PayloadConverter,
272 context: &SerializationContextData,
273 ) -> Result<IncomingError, PayloadConversionError> {
274 Ok(decode_failure(failure, payload_converter, context))
275 }
276}
277
278trait EncodeFailure {
280 fn encode_failure(
281 &self,
282 payload_converter: &PayloadConverter,
283 context: &SerializationContextData,
284 ) -> Result<Failure, PayloadConversionError>;
285}
286
287enum ClassifiedFailure<'a> {
288 Application(&'a ApplicationFailure),
289 ActivityExecution(&'a ActivityExecutionError),
290 ChildWorkflowExecution(&'a ChildWorkflowExecutionError),
291 ChildWorkflowStart(&'a ChildWorkflowStartError),
292 WorkflowSignal(&'a WorkflowSignalError),
293 Generic(&'a (dyn std::error::Error + 'static)),
294}
295
296fn failed_error_conversion_message(
297 original_error: impl std::fmt::Display,
298 converter_error: &PayloadConversionError,
299) -> String {
300 format!(
301 "Failed converting error to failure: {converter_error}, original error message: \
302 {original_error}"
303 )
304}
305
306impl<'a> ClassifiedFailure<'a> {
307 fn from_error(err: &'a (dyn std::error::Error + 'static)) -> Self {
308 if let Some(app) = err.downcast_ref::<ApplicationFailure>() {
309 Self::Application(app)
310 } else if let Some(activity) = err.downcast_ref::<ActivityExecutionError>() {
311 Self::ActivityExecution(activity)
312 } else if let Some(child) = err.downcast_ref::<ChildWorkflowExecutionError>() {
313 Self::ChildWorkflowExecution(child)
314 } else if let Some(child) = err.downcast_ref::<ChildWorkflowStartError>() {
315 Self::ChildWorkflowStart(child)
316 } else if let Some(child_signal) = err.downcast_ref::<WorkflowSignalError>() {
317 Self::WorkflowSignal(child_signal)
318 } else {
319 Self::Generic(err)
320 }
321 }
322
323 fn encode(self) -> Failure {
324 match self {
325 Self::Application(app) => app
326 .encode_failure(
327 &PayloadConverter::default(),
328 &SerializationContextData::None,
329 )
330 .unwrap_or_else(|converter_error| {
331 encode_failed_error_conversion(app, converter_error)
332 }),
333 Self::ActivityExecution(activity) => activity
334 .encode_failure(
335 &PayloadConverter::default(),
336 &SerializationContextData::None,
337 )
338 .unwrap_or_else(|converter_error| {
339 encode_failed_error_conversion(activity, converter_error)
340 }),
341 Self::ChildWorkflowExecution(child) => child
342 .encode_failure(
343 &PayloadConverter::default(),
344 &SerializationContextData::None,
345 )
346 .unwrap_or_else(|converter_error| {
347 encode_failed_error_conversion(child, converter_error)
348 }),
349 Self::ChildWorkflowStart(child) => child
350 .encode_failure(
351 &PayloadConverter::default(),
352 &SerializationContextData::None,
353 )
354 .unwrap_or_else(|converter_error| {
355 encode_failed_error_conversion(child, converter_error)
356 }),
357 Self::WorkflowSignal(signal) => signal
358 .encode_failure(
359 &PayloadConverter::default(),
360 &SerializationContextData::None,
361 )
362 .unwrap_or_else(|converter_error| {
363 encode_failed_error_conversion(signal, converter_error)
364 }),
365 Self::Generic(err) => encode_generic_application_failure(err),
366 }
367 }
368}
369
370impl EncodeFailure for ApplicationFailure {
371 fn encode_failure(
372 &self,
373 payload_converter: &PayloadConverter,
374 context: &SerializationContextData,
375 ) -> Result<Failure, PayloadConversionError> {
376 let details = self
377 .failure_payloads()
378 .map(|details| details.encode(payload_converter, context))
379 .transpose()?;
380 Ok(Failure {
381 message: self.to_string(),
382 cause: self
383 .cause()
384 .map(|cause| Box::new(cause.failure().clone()))
385 .or_else(|| encode_application_failure_cause(self.source_error())),
386 failure_info: Some(FailureInfo::ApplicationFailureInfo(
387 ApplicationFailureInfo {
388 r#type: self.type_name().unwrap_or_default().to_owned(),
389 non_retryable: self.is_non_retryable(),
390 details,
391 next_retry_delay: self.next_retry_delay().and_then(|d| d.try_into().ok()),
392 category: ProtoApplicationErrorCategory::from(self.category()) as i32,
393 },
394 )),
395 ..Default::default()
396 })
397 }
398}
399
400fn encode_application_failure_cause(
401 source: &(dyn std::error::Error + 'static),
402) -> Option<Box<Failure>> {
403 if matches!(
404 ClassifiedFailure::from_error(source),
405 ClassifiedFailure::Application(_) | ClassifiedFailure::Generic(_)
406 ) {
407 source.source().map(encode_error_as_failure).map(Box::new)
408 } else {
409 Some(Box::new(encode_error_as_failure(source)))
410 }
411}
412
413fn encode_error_as_failure(err: &(dyn std::error::Error + 'static)) -> Failure {
414 ClassifiedFailure::from_error(err).encode()
415}
416
417impl EncodeFailure for ActivityExecutionError {
418 fn encode_failure(
419 &self,
420 _: &PayloadConverter,
421 _: &SerializationContextData,
422 ) -> Result<Failure, PayloadConversionError> {
423 Ok(match self {
424 Self::Failed(failure) => failure.failure().clone(),
425 Self::Cancelled(failure) => failure.failure().clone(),
426 Self::Serialization(err) => encode_generic_application_failure(err),
427 })
428 }
429}
430
431impl EncodeFailure for ChildWorkflowExecutionError {
432 fn encode_failure(
433 &self,
434 _: &PayloadConverter,
435 _: &SerializationContextData,
436 ) -> Result<Failure, PayloadConversionError> {
437 Ok(match self {
438 Self::Failed(failure) => failure.failure().clone(),
439 Self::Serialization(_) => encode_generic_application_failure(self),
440 })
441 }
442}
443
444impl EncodeFailure for ChildWorkflowStartError {
445 fn encode_failure(
446 &self,
447 _: &PayloadConverter,
448 _: &SerializationContextData,
449 ) -> Result<Failure, PayloadConversionError> {
450 Ok(match self {
451 Self::Cancelled(failure) => failure.failure().clone(),
452 Self::StartFailed { .. } | Self::Serialization(_) => {
453 encode_generic_application_failure(self)
454 }
455 })
456 }
457}
458
459impl EncodeFailure for WorkflowSignalError {
460 fn encode_failure(
461 &self,
462 _: &PayloadConverter,
463 _: &SerializationContextData,
464 ) -> Result<Failure, PayloadConversionError> {
465 Ok(match self {
466 Self::Failed(failure) => failure.failure().clone(),
467 Self::Serialization(err) => encode_generic_application_failure(err),
468 })
469 }
470}
471
472fn encode_outgoing_activity_error(
473 err: OutgoingActivityError,
474 payload_converter: &PayloadConverter,
475 context: &SerializationContextData,
476) -> Result<Failure, PayloadConversionError> {
477 Ok(match err {
478 OutgoingActivityError::Application(app) => {
479 app.encode_failure(payload_converter, context)?
480 }
481 OutgoingActivityError::Cancelled { details } => Failure {
482 message: "Activity cancelled".to_string(),
483 failure_info: Some(FailureInfo::CanceledFailureInfo(CanceledFailureInfo {
484 details: details
485 .map(|details| details.encode(payload_converter, context))
486 .transpose()?,
487 identity: Default::default(),
488 })),
489 ..Default::default()
490 },
491 })
492}
493
494fn encode_generic_application_failure(err: &(dyn std::error::Error + 'static)) -> Failure {
495 Failure {
496 message: err.to_string(),
497 cause: err.source().map(encode_error_as_failure).map(Box::new),
498 failure_info: Some(FailureInfo::ApplicationFailureInfo(
499 ApplicationFailureInfo::default(),
500 )),
501 ..Default::default()
502 }
503}
504
505fn encode_failed_error_conversion(
506 err: &(dyn std::error::Error + 'static),
507 converter_error: PayloadConversionError,
508) -> Failure {
509 Failure {
510 message: failed_error_conversion_message(err, &converter_error),
511 cause: err.source().map(encode_error_as_failure).map(Box::new),
512 failure_info: Some(FailureInfo::ApplicationFailureInfo(
513 ApplicationFailureInfo::default(),
514 )),
515 ..Default::default()
516 }
517}
518
519fn encode_common_attributes(
520 failure: &mut Failure,
521 payload_converter: &PayloadConverter,
522 context: &SerializationContextData,
523) -> Result<(), PayloadConversionError> {
524 if let Some(cause) = failure.cause.as_deref_mut() {
525 encode_common_attributes(cause, payload_converter, context)?;
526 }
527 failure.encoded_attributes = Some(payload_converter.to_payload(
528 &SerializationContext::new(context, payload_converter),
529 &CommonAttributes {
530 message: std::mem::take(&mut failure.message),
531 stack_trace: std::mem::take(&mut failure.stack_trace),
532 },
533 )?);
534 failure.message = "Encoded failure".to_owned();
535 Ok(())
536}
537
538fn decode_failure(
539 mut failure: Failure,
540 payload_converter: &PayloadConverter,
541 context: &SerializationContextData,
542) -> IncomingError {
543 if let Some(encoded_attributes) = failure.encoded_attributes.clone()
544 && let Ok(attributes) = payload_converter.from_payload::<CommonAttributes>(
545 &SerializationContext::new(context, payload_converter),
546 encoded_attributes,
547 )
548 {
549 failure.message = attributes.message;
550 failure.stack_trace = attributes.stack_trace;
551 }
552 let cause = failure
553 .cause
554 .clone()
555 .map(|cause| decode_failure(*cause, payload_converter, context));
556 match failure.failure_info.clone() {
557 Some(FailureInfo::ApplicationFailureInfo(_)) | None => IncomingError::Application(
558 ApplicationFailure::from_failure(failure, cause, payload_converter, context),
559 ),
560 Some(FailureInfo::TimeoutFailureInfo(failure_info)) => IncomingError::Timeout(
561 TimeoutError::new(failure, failure_info, cause, payload_converter, context),
562 ),
563 Some(FailureInfo::CanceledFailureInfo(failure_info)) => IncomingError::Cancelled(
564 CancelledError::new(failure, failure_info, cause, payload_converter, context),
565 ),
566 Some(FailureInfo::TerminatedFailureInfo(_)) => {
567 IncomingError::Terminated(TerminatedError::new(failure, cause))
568 }
569 Some(FailureInfo::ServerFailureInfo(_)) => {
570 IncomingError::Server(ServerError::new(failure, cause))
571 }
572 Some(FailureInfo::ResetWorkflowFailureInfo(_)) => {
573 IncomingError::ResetWorkflow(ResetWorkflowError::new(failure, cause))
574 }
575 Some(FailureInfo::ActivityFailureInfo(failure_info)) => {
576 IncomingError::Activity(ActivityFailureError::new(failure, failure_info, cause))
577 }
578 Some(FailureInfo::ChildWorkflowExecutionFailureInfo(failure_info)) => {
579 IncomingError::ChildWorkflowExecution(ChildWorkflowFailureError::new(
580 failure,
581 failure_info,
582 cause,
583 ))
584 }
585 Some(FailureInfo::NexusOperationExecutionFailureInfo(_)) => {
586 IncomingError::NexusOperationExecution(IncomingNexusOperationExecutionError::new(
587 failure, cause,
588 ))
589 }
590 Some(FailureInfo::NexusHandlerFailureInfo(_)) => {
591 IncomingError::NexusHandler(IncomingNexusHandlerError::new(failure, cause))
592 }
593 }
594}
595
596#[cfg(test)]
597mod tests {
598 use super::*;
599 use crate::{
600 data_converters::{
601 ActivitySerializationContext, GenericPayloadConverter, SerializationContext,
602 WorkflowSerializationContext,
603 },
604 error::ApplicationErrorCategory,
605 protos::temporal::api::{
606 common::v1::{Payload, Payloads},
607 failure::v1::{
608 ActivityFailureInfo, ChildWorkflowExecutionFailureInfo, NexusHandlerFailureInfo,
609 NexusOperationFailureInfo, ResetWorkflowFailureInfo, ServerFailureInfo,
610 TerminatedFailureInfo, TimeoutFailureInfo, failure::FailureInfo,
611 },
612 },
613 };
614 use rstest::rstest;
615 use std::fmt;
616
617 #[derive(Debug, Clone, Copy)]
618 enum IncomingKind {
619 Application,
620 Timeout,
621 Cancelled,
622 Terminated,
623 Server,
624 ResetWorkflow,
625 Activity,
626 ChildWorkflowExecution,
627 NexusOperationExecution,
628 NexusHandler,
629 }
630
631 #[derive(Debug, Clone, Copy)]
632 enum ActivityExecutionKind {
633 Failed,
634 Cancelled,
635 }
636
637 #[derive(Debug)]
638 struct TestError {
639 message: &'static str,
640 source: Option<Box<dyn std::error::Error + Send + Sync + 'static>>,
641 }
642
643 impl TestError {
644 fn new(
645 message: &'static str,
646 source: Option<Box<dyn std::error::Error + Send + Sync + 'static>>,
647 ) -> Self {
648 Self { message, source }
649 }
650 }
651
652 impl fmt::Display for TestError {
653 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
654 write!(f, "{}", self.message)
655 }
656 }
657
658 impl std::error::Error for TestError {
659 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
660 self.source
661 .as_deref()
662 .map(|source| source as &(dyn std::error::Error + 'static))
663 }
664 }
665
666 struct AlwaysFailsSerialize;
667
668 impl serde::Serialize for AlwaysFailsSerialize {
669 fn serialize<S: serde::Serializer>(&self, _serializer: S) -> Result<S::Ok, S::Error> {
670 Err(serde::ser::Error::custom("serialize boom"))
671 }
672 }
673
674 fn assert_incoming_kind(decoded: &IncomingError, expected: IncomingKind) {
675 match expected {
676 IncomingKind::Application => assert!(matches!(decoded, IncomingError::Application(_))),
677 IncomingKind::Timeout => assert!(matches!(decoded, IncomingError::Timeout(_))),
678 IncomingKind::Cancelled => assert!(matches!(decoded, IncomingError::Cancelled(_))),
679 IncomingKind::Terminated => assert!(matches!(decoded, IncomingError::Terminated(_))),
680 IncomingKind::Server => assert!(matches!(decoded, IncomingError::Server(_))),
681 IncomingKind::ResetWorkflow => {
682 assert!(matches!(decoded, IncomingError::ResetWorkflow(_)))
683 }
684 IncomingKind::Activity => assert!(matches!(decoded, IncomingError::Activity(_))),
685 IncomingKind::ChildWorkflowExecution => {
686 assert!(matches!(decoded, IncomingError::ChildWorkflowExecution(_)))
687 }
688 IncomingKind::NexusOperationExecution => {
689 assert!(matches!(decoded, IncomingError::NexusOperationExecution(_)))
690 }
691 IncomingKind::NexusHandler => {
692 assert!(matches!(decoded, IncomingError::NexusHandler(_)))
693 }
694 }
695 }
696
697 fn convert(err: OutgoingWorkflowError) -> Failure {
698 DefaultFailureConverter::default().to_failure(
699 OutgoingError::Workflow(err),
700 &PayloadConverter::default(),
701 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
702 )
703 }
704
705 fn data_converter() -> crate::data_converters::DataConverter {
706 crate::data_converters::DataConverter::new(
707 PayloadConverter::default(),
708 DefaultFailureConverter::default(),
709 crate::data_converters::DefaultPayloadCodec,
710 )
711 }
712
713 fn cancelled_failure(message: &str) -> Failure {
714 Failure {
715 message: message.to_owned(),
716 failure_info: Some(FailureInfo::CanceledFailureInfo(
717 CanceledFailureInfo::default(),
718 )),
719 ..Default::default()
720 }
721 }
722
723 fn timeout_failure(message: &str) -> Failure {
724 Failure {
725 message: message.to_owned(),
726 failure_info: Some(FailureInfo::TimeoutFailureInfo(
727 TimeoutFailureInfo::default(),
728 )),
729 ..Default::default()
730 }
731 }
732
733 #[test]
734 fn application_failures_preserve_metadata() {
735 let failure = convert(OutgoingWorkflowError::Application(Box::new(
736 ApplicationFailure::builder(anyhow::anyhow!("app boom"))
737 .type_name("MyType".to_owned())
738 .non_retryable(true)
739 .category(ApplicationErrorCategory::Benign)
740 .details(crate::data_converters::RawValue::new(vec![Payload {
741 data: b"details".to_vec(),
742 ..Default::default()
743 }]))
744 .build(),
745 )));
746 let Some(FailureInfo::ApplicationFailureInfo(info)) = failure.failure_info else {
747 panic!("expected application failure info");
748 };
749 assert_eq!(failure.message, "app boom");
750 assert_eq!(info.r#type, "MyType");
751 assert!(info.non_retryable);
752 assert_eq!(info.category(), ProtoApplicationErrorCategory::Benign);
753 assert_eq!(info.details.unwrap().payloads[0].data, b"details".to_vec());
754 }
755
756 #[test]
757 fn application_failures_encode_serializable_details_with_payload_converter() {
758 let failure = convert(OutgoingWorkflowError::Application(Box::new(
759 ApplicationFailure::builder(anyhow::anyhow!("app boom"))
760 .details("detail")
761 .build(),
762 )));
763 let Some(FailureInfo::ApplicationFailureInfo(info)) = failure.failure_info else {
764 panic!("expected application failure info");
765 };
766 let payloads = info.details.expect("details should be present").payloads;
767 let converter = PayloadConverter::default();
768 let details: String = converter
769 .from_payloads(
770 &SerializationContext::new(
771 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
772 &converter,
773 ),
774 payloads,
775 )
776 .unwrap();
777 assert_eq!(details, "detail");
778 }
779
780 #[test]
781 fn application_failures_surface_detail_encoding_errors_with_original_message() {
782 let failure = DefaultFailureConverter::default().to_failure(
783 OutgoingError::Workflow(OutgoingWorkflowError::Application(Box::new(
784 ApplicationFailure::builder(anyhow::anyhow!("app boom"))
785 .details(AlwaysFailsSerialize)
786 .build(),
787 ))),
788 &PayloadConverter::default(),
789 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
790 );
791
792 assert_eq!(
793 failure.message,
794 "Failed converting error to failure: Encoding error: serialize boom, original error message: app boom"
795 );
796 }
797
798 #[test]
799 fn application_failures_decode_details_through_payload_converter() {
800 let converter = PayloadConverter::default();
801 let payloads = converter
802 .to_payloads(
803 &SerializationContext::new(
804 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
805 &converter,
806 ),
807 &"detail",
808 )
809 .unwrap();
810 let failure = Failure {
811 message: "app boom".to_owned(),
812 failure_info: Some(FailureInfo::ApplicationFailureInfo(
813 ApplicationFailureInfo {
814 details: Some(Payloads { payloads }),
815 ..Default::default()
816 },
817 )),
818 ..Default::default()
819 };
820
821 let decoded = DefaultFailureConverter::default()
822 .to_error(
823 failure,
824 &converter,
825 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
826 )
827 .unwrap();
828
829 let IncomingError::Application(app) = decoded else {
830 panic!("expected application error");
831 };
832 assert_eq!(app.details::<String>().unwrap(), Some("detail".to_string()));
833 }
834
835 #[test]
836 fn nested_application_failures_surface_detail_encoding_errors_in_fallback_failure() {
837 let app = ApplicationFailure::new(anyhow::Error::new(TestError::new(
838 "outer wrapper",
839 Some(Box::new(
840 ApplicationFailure::builder(anyhow::anyhow!("inner boom"))
841 .details(AlwaysFailsSerialize)
842 .build(),
843 )),
844 )));
845
846 let converted = convert(OutgoingWorkflowError::Application(Box::new(app)));
847 let cause = converted.cause.expect("expected nested fallback failure");
848 assert_eq!(
849 cause.message,
850 "Failed converting error to failure: Encoding error: serialize boom, original error message: inner boom"
851 );
852 assert!(matches!(
853 cause.failure_info,
854 Some(FailureInfo::ApplicationFailureInfo(_))
855 ));
856 }
857
858 #[test]
859 fn application_failures_do_not_duplicate_their_source_as_cause() {
860 let failure = convert(OutgoingWorkflowError::Application(Box::new(
861 ApplicationFailure::new(anyhow::anyhow!("app boom")),
862 )));
863
864 assert_eq!(failure.message, "app boom");
865 assert!(failure.cause.is_none());
866 }
867
868 #[test]
869 fn application_failures_keep_special_causes_nested() {
870 let activity_failure = Failure {
871 message: "activity failed".to_owned(),
872 failure_info: Some(FailureInfo::ActivityFailureInfo(
873 ActivityFailureInfo::default(),
874 )),
875 ..Default::default()
876 };
877 let app =
878 ApplicationFailure::new(ActivityExecutionError::Failed(ActivityFailureError::new(
879 activity_failure.clone(),
880 ActivityFailureInfo::default(),
881 None,
882 )));
883 let converted = convert(OutgoingWorkflowError::Application(Box::new(app)));
884 assert!(matches!(
885 converted.failure_info,
886 Some(FailureInfo::ApplicationFailureInfo(_))
887 ));
888 assert_eq!(converted.cause.unwrap().as_ref(), &activity_failure);
889 }
890
891 #[test]
892 fn application_failures_fall_back_to_source_error() {
893 let activity_failure = Failure {
894 message: "activity failed".to_owned(),
895 failure_info: Some(FailureInfo::ActivityFailureInfo(
896 ActivityFailureInfo::default(),
897 )),
898 ..Default::default()
899 };
900 let app =
901 ApplicationFailure::new(ActivityExecutionError::Failed(ActivityFailureError::new(
902 activity_failure.clone(),
903 ActivityFailureInfo::default(),
904 None,
905 )));
906
907 assert!(app.cause().is_none());
908
909 let converted = convert(OutgoingWorkflowError::Application(Box::new(app)));
910
911 assert_eq!(converted.cause.unwrap().as_ref(), &activity_failure);
912 }
913
914 #[test]
915 fn application_failures_skip_generic_wrappers_around_known_causes() {
916 let activity_failure = Failure {
917 message: "activity failed".to_owned(),
918 failure_info: Some(FailureInfo::ActivityFailureInfo(
919 ActivityFailureInfo::default(),
920 )),
921 ..Default::default()
922 };
923 let app = ApplicationFailure::new(anyhow::Error::new(TestError::new(
924 "outer wrapper",
925 Some(Box::new(ActivityExecutionError::Failed(
926 ActivityFailureError::new(
927 activity_failure.clone(),
928 ActivityFailureInfo::default(),
929 None,
930 ),
931 ))),
932 )));
933
934 let converted = convert(OutgoingWorkflowError::Application(Box::new(app)));
935
936 assert!(matches!(
937 converted.failure_info,
938 Some(FailureInfo::ApplicationFailureInfo(_))
939 ));
940 assert_eq!(converted.message, "outer wrapper");
941 assert_eq!(converted.cause.unwrap().as_ref(), &activity_failure);
942 }
943
944 #[test]
945 fn application_failures_serialize_unknown_nested_causes_as_application_failures() {
946 let app = ApplicationFailure::new(anyhow::Error::new(TestError::new(
947 "outer wrapper",
948 Some(Box::new(TestError::new("generic inner cause", None))),
949 )));
950
951 let converted = convert(OutgoingWorkflowError::Application(Box::new(app)));
952
953 assert!(matches!(
954 converted.failure_info,
955 Some(FailureInfo::ApplicationFailureInfo(_))
956 ));
957 assert_eq!(converted.message, "outer wrapper",);
958 let cause = converted
959 .cause
960 .clone()
961 .expect("expected nested generic cause");
962 assert_eq!(cause.message, "generic inner cause");
963 assert!(matches!(
964 cause.failure_info,
965 Some(FailureInfo::ApplicationFailureInfo(_))
966 ));
967 assert!(cause.cause.is_none());
968
969 let decoded = DefaultFailureConverter::default()
970 .to_error(
971 converted.clone(),
972 &PayloadConverter::default(),
973 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
974 )
975 .unwrap();
976
977 let IncomingError::Application(decoded_app) = decoded else {
978 panic!("expected application error");
979 };
980 assert_eq!(decoded_app.failure(), Some(&converted));
981 let Some(IncomingError::Application(wrapper)) = decoded_app.cause() else {
982 panic!("expected application cause");
983 };
984 assert_eq!(
985 wrapper.failure().map(|failure| failure.message.as_str()),
986 Some("generic inner cause")
987 );
988 assert!(wrapper.cause().is_none());
989 }
990
991 #[test]
992 fn failure_converter_encodes_and_decodes_cause_chain() {
993 let payload_converter = PayloadConverter::default();
994 let converter = DefaultFailureConverter::new(true);
995 let context = SerializationContextData::Workflow(WorkflowSerializationContext::new());
996 let failure = Failure {
997 message: "outer message".to_owned(),
998 stack_trace: "outer stack trace".to_owned(),
999 cause: Some(Box::new(Failure {
1000 message: "inner message".to_owned(),
1001 stack_trace: "inner stack trace".to_owned(),
1002 failure_info: Some(FailureInfo::ApplicationFailureInfo(
1003 ApplicationFailureInfo::default(),
1004 )),
1005 ..Default::default()
1006 })),
1007 failure_info: Some(FailureInfo::ActivityFailureInfo(
1008 ActivityFailureInfo::default(),
1009 )),
1010 ..Default::default()
1011 };
1012 let activity_error = ActivityExecutionError::Failed(ActivityFailureError::new(
1013 failure,
1014 ActivityFailureInfo::default(),
1015 None,
1016 ));
1017
1018 let failure = converter.to_failure(
1019 OutgoingError::Workflow(OutgoingWorkflowError::ActivityExecution(Box::new(
1020 activity_error,
1021 ))),
1022 &payload_converter,
1023 &context,
1024 );
1025
1026 assert_eq!(failure.message, "Encoded failure");
1027 assert_eq!(failure.cause.as_ref().unwrap().message, "Encoded failure");
1028 assert!(failure.stack_trace.is_empty());
1029 assert!(failure.cause.as_ref().unwrap().stack_trace.is_empty());
1030 let payload_context = SerializationContext::new(&context, &payload_converter);
1031 let outer_attributes: CommonAttributes = payload_converter
1032 .from_payload(
1033 &payload_context,
1034 failure.encoded_attributes.clone().unwrap(),
1035 )
1036 .unwrap();
1037 let inner_attributes: CommonAttributes = payload_converter
1038 .from_payload(
1039 &payload_context,
1040 failure
1041 .cause
1042 .as_ref()
1043 .unwrap()
1044 .encoded_attributes
1045 .clone()
1046 .unwrap(),
1047 )
1048 .unwrap();
1049 assert_eq!(outer_attributes.message, "outer message");
1050 assert_eq!(inner_attributes.message, "inner message");
1051 assert_eq!(outer_attributes.stack_trace, "outer stack trace");
1052 assert_eq!(inner_attributes.stack_trace, "inner stack trace");
1053
1054 let decoded = DefaultFailureConverter::default()
1055 .to_error(failure, &payload_converter, &context)
1056 .unwrap();
1057 assert_eq!(decoded.failure().message, "outer message");
1058 assert_eq!(decoded.failure().stack_trace, "outer stack trace");
1059 let cause = decoded.cause().unwrap().failure();
1060 assert_eq!(cause.message, "inner message");
1061 assert_eq!(cause.stack_trace, "inner stack trace");
1062 }
1063
1064 #[test]
1065 fn start_failed_child_workflow_errors_fall_back_to_application_failures() {
1066 let failure = convert(OutgoingWorkflowError::ChildWorkflowStart(Box::new(
1067 ChildWorkflowStartError::StartFailed {
1068 workflow_id: "wf-id".to_owned(),
1069 workflow_type: "wf-type".to_owned(),
1070 cause: crate::protos::coresdk::child_workflow::StartChildWorkflowExecutionFailedCause::WorkflowAlreadyExists,
1071 },
1072 )));
1073 assert!(matches!(
1074 failure.failure_info,
1075 Some(FailureInfo::ApplicationFailureInfo(_))
1076 ));
1077 assert!(failure.message.contains("Child workflow start failed"));
1078 }
1079
1080 #[test]
1081 fn application_failures_decode_with_metadata_and_proto() {
1082 let failure = Failure {
1083 message: "app boom".to_owned(),
1084 failure_info: Some(FailureInfo::ApplicationFailureInfo(
1085 ApplicationFailureInfo {
1086 r#type: "MyType".to_owned(),
1087 non_retryable: true,
1088 ..Default::default()
1089 },
1090 )),
1091 ..Default::default()
1092 };
1093
1094 let decoded = DefaultFailureConverter::default()
1095 .to_error(
1096 failure.clone(),
1097 &PayloadConverter::default(),
1098 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1099 )
1100 .unwrap();
1101
1102 let IncomingError::Application(app) = decoded else {
1103 panic!("expected application error");
1104 };
1105 assert_eq!(app.type_name(), Some("MyType"));
1106 assert!(app.is_non_retryable());
1107 assert_eq!(app.failure(), Some(&failure));
1108 }
1109
1110 #[test]
1111 fn application_failures_decode_with_normalized_cause() {
1112 let failure = Failure {
1113 message: "app boom".to_owned(),
1114 cause: Some(Box::new(Failure {
1115 message: "timed out".to_owned(),
1116 failure_info: Some(FailureInfo::TimeoutFailureInfo(
1117 TimeoutFailureInfo::default(),
1118 )),
1119 ..Default::default()
1120 })),
1121 failure_info: Some(FailureInfo::ApplicationFailureInfo(
1122 ApplicationFailureInfo::default(),
1123 )),
1124 ..Default::default()
1125 };
1126
1127 let decoded = DefaultFailureConverter::default()
1128 .to_error(
1129 failure.clone(),
1130 &PayloadConverter::default(),
1131 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1132 )
1133 .unwrap();
1134
1135 let IncomingError::Application(app) = decoded else {
1136 panic!("expected application error");
1137 };
1138 assert_eq!(app.failure(), Some(&failure));
1139 assert!(matches!(app.cause(), Some(IncomingError::Timeout(_))));
1140 }
1141
1142 #[test]
1143 fn decoded_application_failures_preserve_cause() {
1144 let failure = Failure {
1145 message: "app boom".to_owned(),
1146 cause: Some(Box::new(timeout_failure("timed out"))),
1147 failure_info: Some(FailureInfo::ApplicationFailureInfo(
1148 ApplicationFailureInfo::default(),
1149 )),
1150 ..Default::default()
1151 };
1152
1153 let decoded = DefaultFailureConverter::default()
1154 .to_error(
1155 failure.clone(),
1156 &PayloadConverter::default(),
1157 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1158 )
1159 .unwrap();
1160
1161 let IncomingError::Application(app) = decoded else {
1162 panic!("expected application error");
1163 };
1164 assert!(app.as_timeout().is_some());
1165
1166 let reencoded = convert(OutgoingWorkflowError::Application(Box::new(app)));
1167
1168 assert_eq!(reencoded.message, failure.message);
1169 assert_eq!(reencoded.cause.as_deref(), failure.cause.as_deref());
1170
1171 let decoded_reencoded = DefaultFailureConverter::default()
1172 .to_error(
1173 reencoded,
1174 &PayloadConverter::default(),
1175 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1176 )
1177 .unwrap();
1178 let IncomingError::Application(roundtripped) = decoded_reencoded else {
1179 panic!("expected application error");
1180 };
1181 assert!(roundtripped.as_timeout().is_some());
1182 }
1183
1184 #[test]
1185 fn application_failures_decode_wrapped_known_causes_without_collapsing_wrapper() {
1186 let failure = Failure {
1187 message: "app boom".to_owned(),
1188 cause: Some(Box::new(Failure {
1189 message: "activity failed".to_owned(),
1190 failure_info: Some(FailureInfo::ActivityFailureInfo(
1191 ActivityFailureInfo::default(),
1192 )),
1193 ..Default::default()
1194 })),
1195 failure_info: Some(FailureInfo::ApplicationFailureInfo(
1196 ApplicationFailureInfo::default(),
1197 )),
1198 ..Default::default()
1199 };
1200
1201 let decoded = DefaultFailureConverter::default()
1202 .to_error(
1203 failure.clone(),
1204 &PayloadConverter::default(),
1205 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1206 )
1207 .unwrap();
1208
1209 let IncomingError::Application(app) = decoded else {
1210 panic!("expected application error");
1211 };
1212 assert_eq!(app.failure(), Some(&failure));
1213 let Some(IncomingError::Activity(activity)) = app.cause() else {
1214 panic!("expected activity cause");
1215 };
1216 assert_eq!(activity.failure().message, "activity failed");
1217 }
1218
1219 #[rstest]
1220 #[case(
1221 FailureInfo::ApplicationFailureInfo(ApplicationFailureInfo::default()),
1222 IncomingKind::Application
1223 )]
1224 #[case(
1225 FailureInfo::TimeoutFailureInfo(TimeoutFailureInfo::default()),
1226 IncomingKind::Timeout
1227 )]
1228 #[case(
1229 FailureInfo::CanceledFailureInfo(CanceledFailureInfo::default()),
1230 IncomingKind::Cancelled
1231 )]
1232 #[case(
1233 FailureInfo::TerminatedFailureInfo(TerminatedFailureInfo::default()),
1234 IncomingKind::Terminated
1235 )]
1236 #[case(
1237 FailureInfo::ServerFailureInfo(ServerFailureInfo::default()),
1238 IncomingKind::Server
1239 )]
1240 #[case(
1241 FailureInfo::ResetWorkflowFailureInfo(ResetWorkflowFailureInfo::default()),
1242 IncomingKind::ResetWorkflow
1243 )]
1244 #[case(
1245 FailureInfo::ActivityFailureInfo(ActivityFailureInfo::default()),
1246 IncomingKind::Activity
1247 )]
1248 #[case(
1249 FailureInfo::ChildWorkflowExecutionFailureInfo(
1250 ChildWorkflowExecutionFailureInfo::default()
1251 ),
1252 IncomingKind::ChildWorkflowExecution
1253 )]
1254 #[case(
1255 FailureInfo::NexusOperationExecutionFailureInfo(NexusOperationFailureInfo::default()),
1256 IncomingKind::NexusOperationExecution
1257 )]
1258 #[case(
1259 FailureInfo::NexusHandlerFailureInfo(NexusHandlerFailureInfo::default()),
1260 IncomingKind::NexusHandler
1261 )]
1262 fn failure_info_decodes_to_expected_incoming_error(
1263 #[case] failure_info: FailureInfo,
1264 #[case] expected: IncomingKind,
1265 ) {
1266 let failure = Failure {
1267 message: "boom".to_owned(),
1268 failure_info: Some(failure_info),
1269 ..Default::default()
1270 };
1271
1272 let decoded = DefaultFailureConverter::default()
1273 .to_error(
1274 failure.clone(),
1275 &PayloadConverter::default(),
1276 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1277 )
1278 .unwrap();
1279
1280 assert_incoming_kind(&decoded, expected);
1281 assert_eq!(decoded.failure(), &failure);
1282 }
1283
1284 #[test]
1285 fn activity_decode_hint_preserves_timeout_reason() {
1286 let failure = Failure {
1287 message: "activity failed".to_owned(),
1288 cause: Some(Box::new(Failure {
1289 message: "timed out".to_owned(),
1290 failure_info: Some(FailureInfo::TimeoutFailureInfo(
1291 TimeoutFailureInfo::default(),
1292 )),
1293 ..Default::default()
1294 })),
1295 failure_info: Some(FailureInfo::ActivityFailureInfo(ActivityFailureInfo {
1296 activity_id: "act-1".to_owned(),
1297 activity_type: Some(crate::protos::temporal::api::common::v1::ActivityType {
1298 name: "test-activity".to_owned(),
1299 }),
1300 scheduled_event_id: 5,
1301 started_event_id: 6,
1302 identity: "worker-1".to_owned(),
1303 retry_state: crate::protos::temporal::api::enums::v1::RetryState::Timeout.into(),
1304 })),
1305 ..Default::default()
1306 };
1307 let data_converter = crate::data_converters::DataConverter::new(
1308 PayloadConverter::default(),
1309 DefaultFailureConverter::default(),
1310 crate::data_converters::DefaultPayloadCodec,
1311 );
1312
1313 let decoded = data_converter
1314 .to_error(
1315 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1316 failure.clone(),
1317 ActivityExecutionDecodeHint { cancelled: false },
1318 )
1319 .unwrap();
1320
1321 let ActivityExecutionError::Failed(decoded_failure) = decoded else {
1322 panic!("expected failed activity execution error");
1323 };
1324 assert_eq!(decoded_failure.failure(), &failure);
1325 assert_eq!(decoded_failure.activity_id(), "act-1");
1326 assert_eq!(decoded_failure.activity_type(), Some("test-activity"));
1327 assert_eq!(decoded_failure.scheduled_event_id(), 5);
1328 assert_eq!(decoded_failure.started_event_id(), 6);
1329 assert_eq!(decoded_failure.identity(), "worker-1");
1330 assert_eq!(
1331 decoded_failure.retry_state(),
1332 crate::error::RetryState::Timeout
1333 );
1334 assert!(matches!(
1335 decoded_failure.cause(),
1336 Some(IncomingError::Timeout(_))
1337 ));
1338 }
1339
1340 #[rstest]
1341 #[case(
1342 cancelled_failure("activity cancelled"),
1343 ActivityExecutionKind::Cancelled,
1344 None
1345 )]
1346 #[case(timeout_failure("timed out"), ActivityExecutionKind::Failed, None)]
1347 #[case(
1348 Failure {
1349 message: "activity task cancelled".to_owned(),
1350 cause: Some(Box::new(cancelled_failure("activity cancelled"))),
1351 failure_info: Some(FailureInfo::ActivityFailureInfo(
1352 ActivityFailureInfo::default(),
1353 )),
1354 ..Default::default()
1355 },
1356 ActivityExecutionKind::Cancelled,
1357 Some(cancelled_failure("activity cancelled"))
1358 )]
1359 fn activity_cancelled_decode_hint_adapts_expected_shape(
1360 #[case] failure: Failure,
1361 #[case] expected_kind: ActivityExecutionKind,
1362 #[case] expected_failure: Option<Failure>,
1363 ) {
1364 let decoded = data_converter()
1365 .to_error(
1366 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1367 failure.clone(),
1368 ActivityExecutionDecodeHint { cancelled: true },
1369 )
1370 .unwrap();
1371
1372 match expected_kind {
1373 ActivityExecutionKind::Failed => {
1374 let ActivityExecutionError::Failed(decoded_failure) = decoded else {
1375 panic!("expected failed activity execution error");
1376 };
1377 assert_eq!(decoded_failure.failure(), &failure);
1378 assert!(decoded_failure.cause().is_none());
1379 }
1380 ActivityExecutionKind::Cancelled => {
1381 let ActivityExecutionError::Cancelled(decoded_failure) = decoded else {
1382 panic!("expected cancelled activity execution error");
1383 };
1384 assert_eq!(
1385 decoded_failure.failure(),
1386 expected_failure.as_ref().unwrap_or(&failure)
1387 );
1388 assert!(decoded_failure.cause().is_none());
1389 }
1390 }
1391 }
1392
1393 #[test]
1394 fn timeout_error_exposes_timeout_info_fields() {
1395 let heartbeat_details = crate::protos::temporal::api::common::v1::Payloads {
1396 payloads: vec![Payload {
1397 data: b"hb".to_vec(),
1398 ..Default::default()
1399 }],
1400 };
1401 let failure = Failure {
1402 message: "timed out".to_owned(),
1403 failure_info: Some(FailureInfo::TimeoutFailureInfo(TimeoutFailureInfo {
1404 timeout_type: crate::protos::temporal::api::enums::v1::TimeoutType::Heartbeat
1405 .into(),
1406 last_heartbeat_details: Some(heartbeat_details.clone()),
1407 })),
1408 ..Default::default()
1409 };
1410
1411 let decoded = DefaultFailureConverter::default()
1412 .to_error(
1413 failure.clone(),
1414 &PayloadConverter::default(),
1415 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1416 )
1417 .unwrap();
1418
1419 let IncomingError::Timeout(timeout) = decoded else {
1420 panic!("expected timeout error");
1421 };
1422 assert_eq!(timeout.timeout_type(), crate::error::TimeoutType::Heartbeat);
1423 assert_eq!(
1424 timeout.raw_last_heartbeat_details(),
1425 Some(heartbeat_details.payloads.as_slice())
1426 );
1427 assert_eq!(timeout.failure(), &failure);
1428 }
1429
1430 #[test]
1431 fn cancelled_error_exposes_details() {
1432 let details = crate::protos::temporal::api::common::v1::Payloads {
1433 payloads: vec![Payload {
1434 data: b"cancel".to_vec(),
1435 ..Default::default()
1436 }],
1437 };
1438 let failure = Failure {
1439 message: "cancelled".to_owned(),
1440 failure_info: Some(FailureInfo::CanceledFailureInfo(CanceledFailureInfo {
1441 details: Some(details.clone()),
1442 identity: Default::default(),
1443 })),
1444 ..Default::default()
1445 };
1446
1447 let decoded = DefaultFailureConverter::default()
1448 .to_error(
1449 failure.clone(),
1450 &PayloadConverter::default(),
1451 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1452 )
1453 .unwrap();
1454
1455 let IncomingError::Cancelled(cancelled) = decoded else {
1456 panic!("expected cancelled error");
1457 };
1458 assert_eq!(cancelled.raw_details(), Some(details.payloads.as_slice()));
1459 assert_eq!(cancelled.failure(), &failure);
1460 }
1461
1462 #[test]
1463 fn child_workflow_decode_hint_preserves_child_failure_proto() {
1464 let failure = Failure {
1465 message: "child workflow failed".to_owned(),
1466 failure_info: Some(FailureInfo::ChildWorkflowExecutionFailureInfo(
1467 ChildWorkflowExecutionFailureInfo {
1468 namespace: "default".to_owned(),
1469 workflow_execution: Some(
1470 crate::protos::temporal::api::common::v1::WorkflowExecution {
1471 workflow_id: "child-id".to_owned(),
1472 run_id: "run-id".to_owned(),
1473 },
1474 ),
1475 workflow_type: Some(crate::protos::temporal::api::common::v1::WorkflowType {
1476 name: "child-type".to_owned(),
1477 }),
1478 initiated_event_id: 11,
1479 started_event_id: 22,
1480 retry_state: crate::protos::temporal::api::enums::v1::RetryState::Timeout
1481 .into(),
1482 },
1483 )),
1484 ..Default::default()
1485 };
1486 let decoded = data_converter()
1487 .to_error(
1488 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1489 failure.clone(),
1490 ChildWorkflowExecutionDecodeHint,
1491 )
1492 .unwrap();
1493
1494 let ChildWorkflowExecutionError::Failed(decoded_failure) = decoded else {
1495 panic!("expected failed child-workflow execution error");
1496 };
1497 assert_eq!(decoded_failure.failure(), &failure);
1498 assert_eq!(decoded_failure.namespace(), "default");
1499 assert_eq!(
1500 decoded_failure
1501 .workflow_execution()
1502 .map(|wf| wf.workflow_id()),
1503 Some("child-id")
1504 );
1505 assert_eq!(
1506 decoded_failure.workflow_execution().map(|wf| wf.run_id()),
1507 Some("run-id")
1508 );
1509 assert_eq!(decoded_failure.workflow_type(), Some("child-type"));
1510 assert_eq!(decoded_failure.initiated_event_id(), 11);
1511 assert_eq!(decoded_failure.started_event_id(), 22);
1512 assert_eq!(
1513 decoded_failure.retry_state(),
1514 crate::error::RetryState::Timeout
1515 );
1516 }
1517
1518 #[rstest]
1519 #[case(
1520 Failure {
1521 message: "child workflow cancelled".to_owned(),
1522 cause: Some(Box::new(cancelled_failure("child workflow cancelled"))),
1523 failure_info: Some(FailureInfo::ChildWorkflowExecutionFailureInfo(
1524 ChildWorkflowExecutionFailureInfo::default(),
1525 )),
1526 ..Default::default()
1527 },
1528 Some(IncomingKind::Cancelled)
1529 )]
1530 #[case(timeout_failure("timed out"), None)]
1531 fn child_workflow_execution_decode_hint_adapts_expected_cause(
1532 #[case] failure: Failure,
1533 #[case] expected_cause: Option<IncomingKind>,
1534 ) {
1535 let decoded = data_converter()
1536 .to_error(
1537 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1538 failure.clone(),
1539 ChildWorkflowExecutionDecodeHint,
1540 )
1541 .unwrap();
1542
1543 let ChildWorkflowExecutionError::Failed(decoded_failure) = decoded else {
1544 panic!("expected failed child-workflow execution error");
1545 };
1546 assert_eq!(decoded_failure.failure(), &failure);
1547 match expected_cause {
1548 Some(expected) => {
1549 let cause = decoded_failure
1550 .cause()
1551 .expect("expected child failure cause");
1552 assert_incoming_kind(cause, expected);
1553 }
1554 None => assert!(decoded_failure.cause().is_none()),
1555 }
1556 }
1557
1558 #[test]
1559 fn child_workflow_start_decode_hint_preserves_top_level_cancellation() {
1560 let failure = Failure {
1561 message: "child start cancelled".to_owned(),
1562 failure_info: Some(FailureInfo::CanceledFailureInfo(
1563 CanceledFailureInfo::default(),
1564 )),
1565 ..Default::default()
1566 };
1567 let decoded = data_converter()
1568 .to_error(
1569 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1570 failure.clone(),
1571 ChildWorkflowStartDecodeHint,
1572 )
1573 .unwrap();
1574
1575 let ChildWorkflowStartError::Cancelled(decoded_failure) = decoded else {
1576 panic!("expected cancelled child-workflow start error");
1577 };
1578 assert_eq!(decoded_failure.failure(), &failure);
1579 assert!(decoded_failure.cause().is_none());
1580 }
1581
1582 #[test]
1583 fn child_workflow_signal_decode_hint_preserves_failure_proto() {
1584 let failure = Failure {
1585 message: "child workflow signal failed".to_owned(),
1586 cause: Some(Box::new(Failure {
1587 message: "timed out".to_owned(),
1588 failure_info: Some(FailureInfo::TimeoutFailureInfo(
1589 TimeoutFailureInfo::default(),
1590 )),
1591 ..Default::default()
1592 })),
1593 ..Default::default()
1594 };
1595 let decoded = data_converter()
1596 .to_error(
1597 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1598 failure.clone(),
1599 WorkflowSignalDecodeHint,
1600 )
1601 .unwrap();
1602
1603 let WorkflowSignalError::Failed(decoded_failure) = decoded else {
1604 panic!("expected failed child-workflow signal error");
1605 };
1606 assert_eq!(decoded_failure.failure(), &failure);
1607 assert!(matches!(
1608 decoded_failure.error(),
1609 IncomingError::Application(_)
1610 ));
1611 assert!(matches!(
1612 decoded_failure.cause(),
1613 Some(IncomingError::Timeout(_))
1614 ));
1615 assert!(std::error::Error::source(&decoded_failure).is_some());
1616 }
1617
1618 #[test]
1619 fn outgoing_cancelled_activity_errors_encode_to_cancelled_failures() {
1620 let failure = DefaultFailureConverter::default().to_failure(
1621 OutgoingError::Activity(OutgoingActivityError::Cancelled { details: None }),
1622 &PayloadConverter::default(),
1623 &SerializationContextData::Activity(ActivitySerializationContext::new()),
1624 );
1625
1626 assert_eq!(failure.message, "Activity cancelled");
1627 assert!(matches!(
1628 failure.failure_info,
1629 Some(FailureInfo::CanceledFailureInfo(_))
1630 ));
1631 }
1632
1633 #[test]
1634 fn outgoing_cancelled_activity_errors_encode_serializable_details_with_payload_converter() {
1635 let failure = DefaultFailureConverter::default().to_failure(
1636 OutgoingError::Activity(OutgoingActivityError::Cancelled {
1637 details: Some("detail".to_string().into()),
1638 }),
1639 &PayloadConverter::default(),
1640 &SerializationContextData::Activity(ActivitySerializationContext::new()),
1641 );
1642
1643 let err = DefaultFailureConverter::default()
1644 .to_error(
1645 failure,
1646 &PayloadConverter::default(),
1647 &SerializationContextData::Activity(ActivitySerializationContext::new()),
1648 )
1649 .unwrap();
1650 let cancelled = err.as_cancelled().unwrap();
1651 let details: String = cancelled.details().unwrap().unwrap();
1652 assert_eq!(details, "detail");
1653 }
1654}