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