1use crate::{
4 WorkflowExecution,
5 data_converters::{
6 DecodablePayloads, GenericPayloadConverter, PayloadConversionError, PayloadConverter,
7 RawValue, SerializationContext, SerializationContextData, TemporalDeserializable,
8 TemporalSerializable,
9 },
10 protos::temporal::api::{
11 common::v1::{Payload, Payloads},
12 enums::v1::{
13 ApplicationErrorCategory as ProtoApplicationErrorCategory,
14 RetryState as ProtoRetryState, TimeoutType as ProtoTimeoutType,
15 },
16 failure::v1::Failure,
17 },
18};
19use std::time::Duration;
20
21#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
23#[non_exhaustive]
24pub enum StartChildWorkflowExecutionFailedCause {
25 Unspecified,
27 WorkflowAlreadyExists,
29 NamespaceNotFound,
31 InvalidVersioningOverride,
33 Unknown,
35}
36
37#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
39#[non_exhaustive]
40pub enum RetryState {
41 Unspecified,
43 InProgress,
45 NonRetryableFailure,
47 Timeout,
49 MaximumAttemptsReached,
51 RetryPolicyNotSet,
53 InternalServerError,
55 CancelRequested,
57 Unknown,
59}
60
61impl RetryState {
62 fn from_raw(value: i32) -> Self {
63 match ProtoRetryState::try_from(value) {
64 Ok(ProtoRetryState::Unspecified) => Self::Unspecified,
65 Ok(ProtoRetryState::InProgress) => Self::InProgress,
66 Ok(ProtoRetryState::NonRetryableFailure) => Self::NonRetryableFailure,
67 Ok(ProtoRetryState::Timeout) => Self::Timeout,
68 Ok(ProtoRetryState::MaximumAttemptsReached) => Self::MaximumAttemptsReached,
69 Ok(ProtoRetryState::RetryPolicyNotSet) => Self::RetryPolicyNotSet,
70 Ok(ProtoRetryState::InternalServerError) => Self::InternalServerError,
71 Ok(ProtoRetryState::CancelRequested) => Self::CancelRequested,
72 Err(_) => Self::Unknown,
73 }
74 }
75}
76
77impl From<ProtoRetryState> for RetryState {
78 fn from(value: ProtoRetryState) -> Self {
79 Self::from_raw(value as i32)
80 }
81}
82
83#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
85#[non_exhaustive]
86pub enum TimeoutType {
87 Unspecified,
89 StartToClose,
91 ScheduleToStart,
93 ScheduleToClose,
95 Heartbeat,
97 Unknown,
99}
100
101impl TimeoutType {
102 fn from_raw(value: i32) -> Self {
103 match ProtoTimeoutType::try_from(value) {
104 Ok(ProtoTimeoutType::Unspecified) => Self::Unspecified,
105 Ok(ProtoTimeoutType::StartToClose) => Self::StartToClose,
106 Ok(ProtoTimeoutType::ScheduleToStart) => Self::ScheduleToStart,
107 Ok(ProtoTimeoutType::ScheduleToClose) => Self::ScheduleToClose,
108 Ok(ProtoTimeoutType::Heartbeat) => Self::Heartbeat,
109 Err(_) => Self::Unknown,
110 }
111 }
112}
113
114impl From<ProtoTimeoutType> for TimeoutType {
115 fn from(value: ProtoTimeoutType) -> Self {
116 Self::from_raw(value as i32)
117 }
118}
119
120trait SerializableFailurePayload: Send + Sync {
123 fn to_payloads(
124 &self,
125 payload_converter: &PayloadConverter,
126 context: &SerializationContextData,
127 ) -> Result<Vec<Payload>, PayloadConversionError>;
128}
129
130impl<T> SerializableFailurePayload for T
131where
132 T: TemporalSerializable + Send + Sync + 'static,
133{
134 fn to_payloads(
135 &self,
136 payload_converter: &PayloadConverter,
137 context: &SerializationContextData,
138 ) -> Result<Vec<Payload>, PayloadConversionError> {
139 payload_converter.to_payloads(&SerializationContext::new(context, payload_converter), self)
140 }
141}
142
143#[derive(derive_more::Debug)]
145pub struct FailurePayloads {
146 repr: FailurePayloadsRepr,
147}
148
149#[derive(derive_more::Debug)]
150enum FailurePayloadsRepr {
151 #[debug("Serializable(...)")]
152 Serializable(#[debug(skip)] Box<dyn SerializableFailurePayload>),
153 Decoded(DecodablePayloads),
154}
155
156impl FailurePayloads {
157 pub(crate) fn encode(
158 &self,
159 payload_converter: &PayloadConverter,
160 context: &SerializationContextData,
161 ) -> Result<Payloads, PayloadConversionError> {
162 let payloads = match &self.repr {
163 FailurePayloadsRepr::Serializable(value) => {
164 value.to_payloads(payload_converter, context)?
165 }
166 FailurePayloadsRepr::Decoded(value) => value.raw().to_vec(),
167 };
168 Ok(Payloads { payloads })
169 }
170
171 pub fn deserialize<T: TemporalDeserializable + 'static>(
173 &self,
174 ) -> Result<T, PayloadConversionError> {
175 match &self.repr {
176 FailurePayloadsRepr::Decoded(value) => value.deserialize(),
177 FailurePayloadsRepr::Serializable(_) => Err(PayloadConversionError::WrongEncoding),
178 }
179 }
180
181 pub fn raw(&self) -> Option<&[Payload]> {
183 match &self.repr {
184 FailurePayloadsRepr::Decoded(value) => Some(value.raw()),
185 FailurePayloadsRepr::Serializable(_) => None,
186 }
187 }
188
189 pub fn into_raw(self) -> Option<RawValue> {
191 match self.repr {
192 FailurePayloadsRepr::Decoded(value) => Some(value.into_raw()),
193 FailurePayloadsRepr::Serializable(_) => None,
194 }
195 }
196}
197
198impl From<DecodablePayloads> for FailurePayloads {
199 fn from(value: DecodablePayloads) -> Self {
200 Self {
201 repr: FailurePayloadsRepr::Decoded(value),
202 }
203 }
204}
205
206impl<T> From<T> for FailurePayloads
207where
208 T: TemporalSerializable + Send + Sync + 'static,
209{
210 fn from(value: T) -> Self {
211 Self {
212 repr: FailurePayloadsRepr::Serializable(Box::new(value)),
213 }
214 }
215}
216
217#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Default)]
219#[non_exhaustive]
220pub enum ApplicationErrorCategory {
221 #[default]
223 Unspecified,
224 Benign,
227}
228
229impl From<ApplicationErrorCategory> for ProtoApplicationErrorCategory {
230 fn from(value: ApplicationErrorCategory) -> Self {
231 match value {
232 ApplicationErrorCategory::Unspecified => ProtoApplicationErrorCategory::Unspecified,
233 ApplicationErrorCategory::Benign => ProtoApplicationErrorCategory::Benign,
234 }
235 }
236}
237
238impl From<ProtoApplicationErrorCategory> for ApplicationErrorCategory {
239 fn from(value: ProtoApplicationErrorCategory) -> Self {
240 match value {
241 ProtoApplicationErrorCategory::Unspecified => ApplicationErrorCategory::Unspecified,
242 ProtoApplicationErrorCategory::Benign => ApplicationErrorCategory::Benign,
243 }
244 }
245}
246
247#[derive(Debug, bon::Builder)]
249#[builder(start_fn = builder, state_mod(vis = "pub"))]
250pub struct ApplicationFailure {
251 #[builder(start_fn, into)]
252 source: Box<dyn std::error::Error + Send + Sync>,
253 type_name: Option<String>,
254 #[builder(default)]
255 non_retryable: bool,
256 next_retry_delay: Option<Duration>,
257 #[builder(default = ApplicationErrorCategory::Unspecified)]
258 category: ApplicationErrorCategory,
259 #[builder(into)]
260 details: Option<FailurePayloads>,
261 failure: Option<Failure>,
262 cause: Option<Box<IncomingError>>,
263}
264
265impl ApplicationFailure {
266 pub fn new(source: impl Into<Box<dyn std::error::Error + Send + Sync>>) -> Self {
268 Self {
269 source: source.into(),
270 type_name: None,
271 non_retryable: false,
272 next_retry_delay: None,
273 category: ApplicationErrorCategory::Unspecified,
274 details: None,
275 failure: None,
276 cause: None,
277 }
278 }
279
280 pub fn non_retryable(source: impl Into<Box<dyn std::error::Error + Send + Sync>>) -> Self {
282 Self {
283 non_retryable: true,
284 ..Self::new(source)
285 }
286 }
287
288 pub fn source_error(&self) -> &(dyn std::error::Error + Send + Sync + 'static) {
290 &*self.source as &(dyn std::error::Error + Send + Sync + 'static)
291 }
292
293 pub fn type_name(&self) -> Option<&str> {
295 self.type_name.as_deref()
296 }
297
298 pub fn is_non_retryable(&self) -> bool {
300 self.non_retryable
301 }
302
303 pub fn next_retry_delay(&self) -> Option<Duration> {
305 self.next_retry_delay
306 }
307
308 pub fn category(&self) -> ApplicationErrorCategory {
310 self.category
311 }
312
313 pub fn details<T: TemporalDeserializable + 'static>(
315 &self,
316 ) -> Result<Option<T>, PayloadConversionError> {
317 self.details
318 .as_ref()
319 .map(FailurePayloads::deserialize)
320 .transpose()
321 }
322
323 pub fn raw_details(&self) -> Option<&[Payload]> {
325 self.details.as_ref().and_then(FailurePayloads::raw)
326 }
327
328 pub(crate) fn failure_payloads(&self) -> Option<&FailurePayloads> {
329 self.details.as_ref()
330 }
331
332 pub fn failure(&self) -> Option<&Failure> {
334 self.failure.as_ref()
335 }
336
337 pub fn into_failure(self) -> Option<Failure> {
339 self.failure
340 }
341
342 pub fn cause(&self) -> Option<&IncomingError> {
344 self.cause.as_deref()
345 }
346
347 pub fn as_timeout(&self) -> Option<&TimeoutError> {
350 self.cause().and_then(IncomingError::as_timeout)
351 }
352
353 pub fn as_cancelled(&self) -> Option<&CancelledError> {
356 self.cause().and_then(IncomingError::as_cancelled)
357 }
358
359 pub(crate) fn from_failure(
360 failure: Failure,
361 cause: Option<IncomingError>,
362 payload_converter: &PayloadConverter,
363 context: &SerializationContextData,
364 ) -> Self {
365 let app_info = failure
366 .maybe_application_failure()
367 .cloned()
368 .unwrap_or_default();
369 let type_name = (!app_info.r#type.is_empty()).then_some(app_info.r#type.clone());
370 Self {
371 source: failure.message.clone().into(),
372 type_name,
373 non_retryable: app_info.non_retryable,
374 next_retry_delay: app_info.next_retry_delay.and_then(|d| d.try_into().ok()),
375 category: app_info.category().into(),
376 details: app_info.details.map(|details| {
377 FailurePayloads::from(DecodablePayloads::new(
378 details.payloads,
379 payload_converter.clone(),
380 context.clone(),
381 ))
382 }),
383 failure: Some(failure),
384 cause: cause.map(Box::new),
385 }
386 }
387}
388
389impl std::fmt::Display for ApplicationFailure {
390 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
391 write!(f, "{}", self.source)
392 }
393}
394
395impl std::error::Error for ApplicationFailure {
396 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
397 self.cause
398 .as_deref()
399 .map(|cause| cause as &(dyn std::error::Error + 'static))
400 .or_else(|| Some(self.source.as_ref()))
401 }
402}
403
404impl From<anyhow::Error> for ApplicationFailure {
405 fn from(value: anyhow::Error) -> Self {
406 Self::new(value)
407 }
408}
409
410impl From<PayloadConversionError> for ApplicationFailure {
411 fn from(value: PayloadConversionError) -> Self {
412 Self::new(value)
413 }
414}
415
416#[derive(Debug, thiserror::Error)]
418#[non_exhaustive]
419pub enum OutgoingError {
420 #[error(transparent)]
422 Activity(#[from] OutgoingActivityError),
423 #[error(transparent)]
425 Workflow(#[from] OutgoingWorkflowError),
426}
427
428#[derive(Debug, thiserror::Error)]
430#[non_exhaustive]
431pub enum OutgoingActivityError {
432 #[error(transparent)]
434 Application(#[from] Box<ApplicationFailure>),
435 #[error("Activity cancelled")]
437 Cancelled {
438 details: Option<FailurePayloads>,
440 },
441}
442
443#[derive(Debug, thiserror::Error)]
445#[non_exhaustive]
446pub enum OutgoingWorkflowError {
447 #[error(transparent)]
449 Application(#[from] Box<ApplicationFailure>),
450 #[error(transparent)]
452 PayloadConversion(#[from] PayloadConversionError),
453 #[error(transparent)]
455 ActivityExecution(#[from] Box<ActivityExecutionError>),
456 #[error(transparent)]
458 ChildWorkflowExecution(#[from] Box<ChildWorkflowExecutionError>),
459 #[error(transparent)]
461 ChildWorkflowStart(#[from] Box<ChildWorkflowStartError>),
462 #[error(transparent)]
464 WorkflowSignal(#[from] Box<WorkflowSignalError>),
465 #[error(transparent)]
467 CancelExternalWorkflow(#[from] Box<CancelExternalWorkflowError>),
468}
469
470impl OutgoingWorkflowError {
471 pub fn as_cancelled(&self) -> Option<&CancelledError> {
474 match self {
475 Self::Application(err) => err.as_cancelled(),
476 Self::PayloadConversion(_) => None,
477 Self::ActivityExecution(err) => err.as_cancelled(),
478 Self::ChildWorkflowExecution(err) => err.as_cancelled(),
479 Self::ChildWorkflowStart(err) => err.as_cancelled(),
480 Self::WorkflowSignal(err) => err.as_cancelled(),
481 Self::CancelExternalWorkflow(err) => err.as_cancelled(),
482 }
483 }
484}
485
486impl From<anyhow::Error> for OutgoingWorkflowError {
487 fn from(value: anyhow::Error) -> Self {
488 Self::Application(Box::new(ApplicationFailure::new(value)))
489 }
490}
491
492impl From<ApplicationFailure> for OutgoingWorkflowError {
493 fn from(value: ApplicationFailure) -> Self {
494 Self::Application(Box::new(value))
495 }
496}
497
498impl From<ActivityExecutionError> for OutgoingWorkflowError {
499 fn from(value: ActivityExecutionError) -> Self {
500 match value {
501 ActivityExecutionError::Serialization(err) => Self::PayloadConversion(err),
502 other => Self::ActivityExecution(Box::new(other)),
503 }
504 }
505}
506
507impl From<ChildWorkflowExecutionError> for OutgoingWorkflowError {
508 fn from(value: ChildWorkflowExecutionError) -> Self {
509 match value {
510 ChildWorkflowExecutionError::Serialization(err) => Self::PayloadConversion(err),
511 other => Self::ChildWorkflowExecution(Box::new(other)),
512 }
513 }
514}
515
516impl From<ChildWorkflowStartError> for OutgoingWorkflowError {
517 fn from(value: ChildWorkflowStartError) -> Self {
518 match value {
519 ChildWorkflowStartError::Serialization(err) => Self::PayloadConversion(err),
520 other => Self::ChildWorkflowStart(Box::new(other)),
521 }
522 }
523}
524
525impl From<WorkflowSignalError> for OutgoingWorkflowError {
526 fn from(value: WorkflowSignalError) -> Self {
527 match value {
528 WorkflowSignalError::Serialization(err) => Self::PayloadConversion(err),
529 other => Self::WorkflowSignal(Box::new(other)),
530 }
531 }
532}
533
534impl From<CancelExternalWorkflowError> for OutgoingWorkflowError {
535 fn from(value: CancelExternalWorkflowError) -> Self {
536 match value {
537 CancelExternalWorkflowError::Serialization(err) => Self::PayloadConversion(err),
538 other => Self::CancelExternalWorkflow(Box::new(other)),
539 }
540 }
541}
542
543#[derive(Debug)]
545#[non_exhaustive]
546pub enum IncomingError {
547 Application(ApplicationFailure),
549 Timeout(TimeoutError),
551 Cancelled(CancelledError),
553 Terminated(TerminatedError),
555 Server(ServerError),
557 ResetWorkflow(ResetWorkflowError),
559 Activity(ActivityFailureError),
561 ChildWorkflowExecution(ChildWorkflowFailureError),
563 NexusOperationExecution(IncomingNexusOperationExecutionError),
565 NexusHandler(IncomingNexusHandlerError),
567}
568
569impl IncomingError {
570 pub fn failure(&self) -> &Failure {
572 match self {
573 IncomingError::Application(err) => err
574 .failure()
575 .expect("decoded application failures retain their original proto"),
576 IncomingError::Timeout(err) => err.failure(),
577 IncomingError::Cancelled(err) => err.failure(),
578 IncomingError::Terminated(err) => err.failure(),
579 IncomingError::Server(err) => err.failure(),
580 IncomingError::ResetWorkflow(err) => err.failure(),
581 IncomingError::Activity(err) => err.failure(),
582 IncomingError::ChildWorkflowExecution(err) => err.failure(),
583 IncomingError::NexusOperationExecution(err) => err.failure(),
584 IncomingError::NexusHandler(err) => err.failure(),
585 }
586 }
587
588 pub fn cause(&self) -> Option<&IncomingError> {
590 match self {
591 IncomingError::Application(err) => err.cause(),
592 IncomingError::Timeout(err) => err.cause(),
593 IncomingError::Cancelled(err) => err.cause(),
594 IncomingError::Terminated(err) => err.cause(),
595 IncomingError::Server(err) => err.cause(),
596 IncomingError::ResetWorkflow(err) => err.cause(),
597 IncomingError::Activity(err) => err.cause(),
598 IncomingError::ChildWorkflowExecution(err) => err.cause(),
599 IncomingError::NexusOperationExecution(err) => err.cause(),
600 IncomingError::NexusHandler(err) => err.cause(),
601 }
602 }
603
604 pub fn into_failure(self) -> Failure {
606 match self {
607 IncomingError::Application(err) => err
608 .into_failure()
609 .expect("decoded application failures retain their original proto"),
610 IncomingError::Timeout(err) => err.into_failure(),
611 IncomingError::Cancelled(err) => err.into_failure(),
612 IncomingError::Terminated(err) => err.into_failure(),
613 IncomingError::Server(err) => err.into_failure(),
614 IncomingError::ResetWorkflow(err) => err.into_failure(),
615 IncomingError::Activity(err) => err.into_failure(),
616 IncomingError::ChildWorkflowExecution(err) => err.into_failure(),
617 IncomingError::NexusOperationExecution(err) => err.into_failure(),
618 IncomingError::NexusHandler(err) => err.into_failure(),
619 }
620 }
621
622 pub fn as_timeout(&self) -> Option<&TimeoutError> {
624 match self {
625 IncomingError::Timeout(err) => Some(err),
626 _ => None,
627 }
628 }
629
630 pub fn as_cancelled(&self) -> Option<&CancelledError> {
632 match self {
633 IncomingError::Cancelled(err) => Some(err),
634 _ => None,
635 }
636 }
637}
638
639impl std::fmt::Display for IncomingError {
640 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
641 match self {
642 IncomingError::Application(err) => err.fmt(f),
643 IncomingError::Timeout(err) => err.fmt(f),
644 IncomingError::Cancelled(err) => err.fmt(f),
645 IncomingError::Terminated(err) => err.fmt(f),
646 IncomingError::Server(err) => err.fmt(f),
647 IncomingError::ResetWorkflow(err) => err.fmt(f),
648 IncomingError::Activity(err) => err.fmt(f),
649 IncomingError::ChildWorkflowExecution(err) => err.fmt(f),
650 IncomingError::NexusOperationExecution(err) => err.fmt(f),
651 IncomingError::NexusHandler(err) => err.fmt(f),
652 }
653 }
654}
655
656impl std::error::Error for IncomingError {
657 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
658 match self {
659 IncomingError::Application(err) => Some(err),
660 IncomingError::Timeout(err) => Some(err),
661 IncomingError::Cancelled(err) => Some(err),
662 IncomingError::Terminated(err) => Some(err),
663 IncomingError::Server(err) => Some(err),
664 IncomingError::ResetWorkflow(err) => Some(err),
665 IncomingError::Activity(err) => Some(err),
666 IncomingError::ChildWorkflowExecution(err) => Some(err),
667 IncomingError::NexusOperationExecution(err) => Some(err),
668 IncomingError::NexusHandler(err) => Some(err),
669 }
670 }
671}
672
673macro_rules! impl_incoming_failure_wrapper {
674 ($name:ident) => {
675 impl $name {
676 pub fn failure(&self) -> &Failure {
678 &self.failure
679 }
680
681 pub fn cause(&self) -> Option<&IncomingError> {
683 self.cause.as_deref()
684 }
685
686 pub fn into_failure(self) -> Failure {
688 self.failure
689 }
690
691 pub fn into_parts(self) -> (Failure, Option<IncomingError>) {
693 (self.failure, self.cause.map(|cause| *cause))
694 }
695 }
696
697 impl std::fmt::Display for $name {
698 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
699 self.failure.fmt(f)
700 }
701 }
702
703 impl std::error::Error for $name {
704 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
705 self.cause
706 .as_deref()
707 .map(|cause| cause as &(dyn std::error::Error + 'static))
708 }
709 }
710 };
711}
712
713macro_rules! incoming_failure_wrapper {
714 ($name:ident, $doc:literal) => {
715 #[doc = $doc]
716 #[derive(Debug)]
717 pub struct $name {
718 failure: Failure,
719 cause: Option<Box<IncomingError>>,
720 }
721
722 impl $name {
723 pub(crate) fn new(failure: Failure, cause: Option<IncomingError>) -> Self {
725 Self {
726 failure,
727 cause: cause.map(Box::new),
728 }
729 }
730 }
731
732 impl_incoming_failure_wrapper!($name);
733 };
734}
735
736#[derive(Debug)]
738pub struct TimeoutError {
739 failure: Failure,
740 cause: Option<Box<IncomingError>>,
741 timeout_type: TimeoutType,
742 last_heartbeat_details: Option<DecodablePayloads>,
743}
744
745impl TimeoutError {
746 pub(crate) fn new(
748 failure: Failure,
749 failure_info: crate::protos::temporal::api::failure::v1::TimeoutFailureInfo,
750 cause: Option<IncomingError>,
751 payload_converter: &PayloadConverter,
752 context: &SerializationContextData,
753 ) -> Self {
754 Self {
755 failure,
756 cause: cause.map(Box::new),
757 timeout_type: TimeoutType::from_raw(failure_info.timeout_type),
758 last_heartbeat_details: failure_info.last_heartbeat_details.map(|details| {
759 DecodablePayloads::new(details.payloads, payload_converter.clone(), context.clone())
760 }),
761 }
762 }
763
764 pub fn timeout_type(&self) -> TimeoutType {
766 self.timeout_type
767 }
768
769 pub fn last_heartbeat_details<T: TemporalDeserializable + 'static>(
771 &self,
772 ) -> Result<Option<T>, PayloadConversionError> {
773 self.last_heartbeat_details
774 .as_ref()
775 .map(DecodablePayloads::deserialize)
776 .transpose()
777 }
778
779 pub fn raw_last_heartbeat_details(&self) -> Option<&[Payload]> {
781 self.last_heartbeat_details
782 .as_ref()
783 .map(DecodablePayloads::raw)
784 }
785}
786
787impl_incoming_failure_wrapper!(TimeoutError);
788
789#[derive(Debug)]
791pub struct CancelledError {
792 failure: Failure,
793 cause: Option<Box<IncomingError>>,
794 details: Option<DecodablePayloads>,
795}
796
797impl CancelledError {
798 pub(crate) fn new(
800 failure: Failure,
801 failure_info: crate::protos::temporal::api::failure::v1::CanceledFailureInfo,
802 cause: Option<IncomingError>,
803 payload_converter: &PayloadConverter,
804 context: &SerializationContextData,
805 ) -> Self {
806 Self {
807 failure,
808 cause: cause.map(Box::new),
809 details: failure_info.details.map(|details| {
810 DecodablePayloads::new(details.payloads, payload_converter.clone(), context.clone())
811 }),
812 }
813 }
814
815 pub fn details<T: TemporalDeserializable + 'static>(
818 &self,
819 ) -> Result<Option<T>, PayloadConversionError> {
820 self.details
821 .as_ref()
822 .map(DecodablePayloads::deserialize)
823 .transpose()
824 }
825
826 pub fn raw_details(&self) -> Option<&[Payload]> {
828 self.details.as_ref().map(DecodablePayloads::raw)
829 }
830}
831
832impl_incoming_failure_wrapper!(CancelledError);
833incoming_failure_wrapper!(TerminatedError, "A normalized terminated failure.");
834incoming_failure_wrapper!(ServerError, "A normalized server failure.");
835incoming_failure_wrapper!(ResetWorkflowError, "A normalized reset-workflow failure.");
836
837#[derive(Debug)]
839pub struct ActivityFailureError {
840 failure: Failure,
841 cause: Option<Box<IncomingError>>,
842 activity_id: String,
843 activity_type: Option<String>,
844 scheduled_event_id: i64,
845 started_event_id: i64,
846 identity: String,
847 retry_state: RetryState,
848}
849
850impl ActivityFailureError {
851 pub(crate) fn new(
853 failure: Failure,
854 failure_info: crate::protos::temporal::api::failure::v1::ActivityFailureInfo,
855 cause: Option<IncomingError>,
856 ) -> Self {
857 let retry_state = RetryState::from_raw(failure_info.retry_state);
858 Self {
859 failure,
860 cause: cause.map(Box::new),
861 activity_id: failure_info.activity_id,
862 activity_type: failure_info
863 .activity_type
864 .map(|activity_type| activity_type.name),
865 scheduled_event_id: failure_info.scheduled_event_id,
866 started_event_id: failure_info.started_event_id,
867 identity: failure_info.identity,
868 retry_state,
869 }
870 }
871
872 pub fn activity_id(&self) -> &str {
874 &self.activity_id
875 }
876
877 pub fn activity_type(&self) -> Option<&str> {
879 self.activity_type.as_deref()
880 }
881
882 pub fn scheduled_event_id(&self) -> i64 {
884 self.scheduled_event_id
885 }
886
887 pub fn started_event_id(&self) -> i64 {
889 self.started_event_id
890 }
891
892 pub fn identity(&self) -> &str {
894 &self.identity
895 }
896
897 pub fn retry_state(&self) -> RetryState {
899 self.retry_state
900 }
901
902 pub fn as_timeout(&self) -> Option<&TimeoutError> {
905 self.cause().and_then(IncomingError::as_timeout)
906 }
907
908 pub fn as_cancelled(&self) -> Option<&CancelledError> {
911 self.cause().and_then(IncomingError::as_cancelled)
912 }
913}
914
915impl_incoming_failure_wrapper!(ActivityFailureError);
916#[derive(Debug)]
918pub struct ChildWorkflowFailureError {
919 failure: Failure,
920 cause: Option<Box<IncomingError>>,
921 namespace: String,
922 workflow_execution: Option<WorkflowExecution>,
923 workflow_type: Option<String>,
924 initiated_event_id: i64,
925 started_event_id: i64,
926 retry_state: RetryState,
927}
928
929impl ChildWorkflowFailureError {
930 pub(crate) fn new(
932 failure: Failure,
933 failure_info: crate::protos::temporal::api::failure::v1::ChildWorkflowExecutionFailureInfo,
934 cause: Option<IncomingError>,
935 ) -> Self {
936 let retry_state = RetryState::from_raw(failure_info.retry_state);
937 Self {
938 failure,
939 cause: cause.map(Box::new),
940 namespace: failure_info.namespace,
941 workflow_execution: failure_info.workflow_execution.map(Into::into),
942 workflow_type: failure_info
943 .workflow_type
944 .map(|workflow_type| workflow_type.name),
945 initiated_event_id: failure_info.initiated_event_id,
946 started_event_id: failure_info.started_event_id,
947 retry_state,
948 }
949 }
950
951 pub fn namespace(&self) -> &str {
953 &self.namespace
954 }
955
956 pub fn workflow_execution(&self) -> Option<&WorkflowExecution> {
958 self.workflow_execution.as_ref()
959 }
960
961 pub fn workflow_type(&self) -> Option<&str> {
963 self.workflow_type.as_deref()
964 }
965
966 pub fn initiated_event_id(&self) -> i64 {
968 self.initiated_event_id
969 }
970
971 pub fn started_event_id(&self) -> i64 {
973 self.started_event_id
974 }
975
976 pub fn retry_state(&self) -> RetryState {
978 self.retry_state
979 }
980
981 pub fn as_timeout(&self) -> Option<&TimeoutError> {
984 self.cause().and_then(IncomingError::as_timeout)
985 }
986
987 pub fn as_cancelled(&self) -> Option<&CancelledError> {
990 self.cause().and_then(IncomingError::as_cancelled)
991 }
992}
993
994impl_incoming_failure_wrapper!(ChildWorkflowFailureError);
995incoming_failure_wrapper!(
996 IncomingNexusOperationExecutionError,
997 "A normalized nexus operation failure wrapper."
998);
999incoming_failure_wrapper!(
1000 IncomingNexusHandlerError,
1001 "A normalized nexus handler failure wrapper."
1002);
1003
1004#[derive(Debug, thiserror::Error)]
1006#[non_exhaustive]
1007pub enum ActivityExecutionError {
1008 #[error("Activity failed: {}", .0.failure().message)]
1010 Failed(#[source] ActivityFailureError),
1011 #[error("Activity cancelled: {}", .0.failure().message)]
1013 Cancelled(#[source] CancelledError),
1014 #[error("Payload conversion failed: {0}")]
1016 Serialization(#[from] PayloadConversionError),
1017}
1018
1019impl ActivityExecutionError {
1020 pub fn failure(&self) -> Option<&Failure> {
1022 match self {
1023 ActivityExecutionError::Failed(err) => Some(err.failure()),
1024 ActivityExecutionError::Cancelled(err) => Some(err.failure()),
1025 ActivityExecutionError::Serialization(_) => None,
1026 }
1027 }
1028
1029 pub fn cause(&self) -> Option<&IncomingError> {
1031 match self {
1032 ActivityExecutionError::Failed(err) => err.cause(),
1033 ActivityExecutionError::Cancelled(err) => err.cause(),
1034 ActivityExecutionError::Serialization(_) => None,
1035 }
1036 }
1037
1038 pub fn reason(&self) -> Option<&IncomingError> {
1040 match self {
1041 ActivityExecutionError::Failed(err) => err.cause(),
1042 ActivityExecutionError::Cancelled(_) | ActivityExecutionError::Serialization(_) => None,
1043 }
1044 }
1045
1046 pub fn as_timeout(&self) -> Option<&TimeoutError> {
1049 match self {
1050 ActivityExecutionError::Failed(err) => err.as_timeout(),
1051 ActivityExecutionError::Serialization(_) | ActivityExecutionError::Cancelled(_) => None,
1052 }
1053 }
1054
1055 pub fn as_cancelled(&self) -> Option<&CancelledError> {
1058 match self {
1059 ActivityExecutionError::Failed(err) => err.as_cancelled(),
1060 ActivityExecutionError::Cancelled(err) => Some(err),
1061 ActivityExecutionError::Serialization(_) => None,
1062 }
1063 }
1064}
1065
1066#[derive(Debug, thiserror::Error)]
1068#[non_exhaustive]
1069pub enum ChildWorkflowStartError {
1070 #[error("Child workflow start cancelled: {}", .0.failure().message)]
1072 Cancelled(#[source] Box<CancelledError>),
1073 #[error(
1075 "Child workflow start failed: workflow_id={workflow_id}, workflow_type={workflow_type}, cause={cause:?}"
1076 )]
1077 StartFailed {
1078 workflow_id: String,
1080 workflow_type: String,
1082 cause: StartChildWorkflowExecutionFailedCause,
1084 },
1085 #[error("Payload conversion failed: {0}")]
1087 Serialization(#[from] PayloadConversionError),
1088}
1089
1090impl ChildWorkflowStartError {
1091 pub fn failure(&self) -> Option<&Failure> {
1093 match self {
1094 ChildWorkflowStartError::Cancelled(err) => Some(err.failure()),
1095 ChildWorkflowStartError::StartFailed { .. }
1096 | ChildWorkflowStartError::Serialization(_) => None,
1097 }
1098 }
1099
1100 pub fn cause(&self) -> Option<&IncomingError> {
1102 match self {
1103 ChildWorkflowStartError::Cancelled(err) => err.cause(),
1104 ChildWorkflowStartError::StartFailed { .. }
1105 | ChildWorkflowStartError::Serialization(_) => None,
1106 }
1107 }
1108
1109 pub fn as_cancelled(&self) -> Option<&CancelledError> {
1112 match self {
1113 ChildWorkflowStartError::Cancelled(err) => Some(err),
1114 ChildWorkflowStartError::StartFailed { .. }
1115 | ChildWorkflowStartError::Serialization(_) => None,
1116 }
1117 }
1118}
1119
1120#[derive(Debug, thiserror::Error)]
1122#[non_exhaustive]
1123pub enum ChildWorkflowExecutionError {
1124 #[error("Child workflow failed: {}", .0.failure().message)]
1126 Failed(#[source] Box<ChildWorkflowFailureError>),
1127 #[error("Payload conversion failed: {0}")]
1129 Serialization(#[from] PayloadConversionError),
1130}
1131
1132impl ChildWorkflowExecutionError {
1133 pub fn failure(&self) -> Option<&Failure> {
1135 match self {
1136 ChildWorkflowExecutionError::Failed(err) => Some(err.failure()),
1137 ChildWorkflowExecutionError::Serialization(_) => None,
1138 }
1139 }
1140
1141 pub fn cause(&self) -> Option<&IncomingError> {
1143 match self {
1144 ChildWorkflowExecutionError::Failed(err) => err.cause(),
1145 ChildWorkflowExecutionError::Serialization(_) => None,
1146 }
1147 }
1148
1149 pub fn reason(&self) -> Option<&IncomingError> {
1151 match self {
1152 ChildWorkflowExecutionError::Failed(err) => err.cause(),
1153 ChildWorkflowExecutionError::Serialization(_) => None,
1154 }
1155 }
1156
1157 pub fn as_timeout(&self) -> Option<&TimeoutError> {
1160 match self {
1161 ChildWorkflowExecutionError::Failed(err) => err.as_timeout(),
1162 ChildWorkflowExecutionError::Serialization(_) => None,
1163 }
1164 }
1165
1166 pub fn as_cancelled(&self) -> Option<&CancelledError> {
1169 match self {
1170 ChildWorkflowExecutionError::Failed(err) => err.as_cancelled(),
1171 ChildWorkflowExecutionError::Serialization(_) => None,
1172 }
1173 }
1174}
1175
1176#[derive(Debug, thiserror::Error)]
1178#[non_exhaustive]
1179pub enum WorkflowSignalError {
1180 #[error("Workflow not found: {}", .0.failure().message)]
1182 NotFound(#[source] Box<WorkflowSignalFailureError>),
1183 #[error("Child workflow signal failed: {}", .0.failure().message)]
1185 Failed(#[source] Box<WorkflowSignalFailureError>),
1186 #[error("Signal payload conversion failed: {0}")]
1188 Serialization(#[from] PayloadConversionError),
1189}
1190
1191#[derive(Debug, thiserror::Error)]
1193pub enum CancelExternalWorkflowError {
1194 #[error("Workflow not found: {}", .0.failure().message)]
1196 NotFound(#[source] Box<WorkflowCancelFailureError>),
1197 #[error("External workflow cancellation request failed: {}", .0.failure().message)]
1199 Failed(#[source] Box<WorkflowCancelFailureError>),
1200 #[error("External workflow cancellation failure conversion failed: {0}")]
1202 Serialization(#[from] PayloadConversionError),
1203}
1204
1205impl CancelExternalWorkflowError {
1206 pub fn failure(&self) -> Option<&Failure> {
1208 match self {
1209 Self::NotFound(err) | Self::Failed(err) => Some(err.failure()),
1210 Self::Serialization(_) => None,
1211 }
1212 }
1213
1214 pub fn cause(&self) -> Option<&IncomingError> {
1216 match self {
1217 Self::NotFound(err) | Self::Failed(err) => err.cause(),
1218 Self::Serialization(_) => None,
1219 }
1220 }
1221
1222 pub fn reason(&self) -> Option<&IncomingError> {
1224 match self {
1225 Self::NotFound(err) | Self::Failed(err) => Some(err.error()),
1226 Self::Serialization(_) => None,
1227 }
1228 }
1229
1230 pub fn as_cancelled(&self) -> Option<&CancelledError> {
1232 self.reason()?.as_cancelled()
1233 }
1234}
1235
1236impl WorkflowSignalError {
1237 pub fn failure(&self) -> Option<&Failure> {
1239 match self {
1240 WorkflowSignalError::NotFound(err) | WorkflowSignalError::Failed(err) => {
1241 Some(err.failure())
1242 }
1243 WorkflowSignalError::Serialization(_) => None,
1244 }
1245 }
1246
1247 pub fn cause(&self) -> Option<&IncomingError> {
1249 match self {
1250 WorkflowSignalError::NotFound(err) | WorkflowSignalError::Failed(err) => err.cause(),
1251 WorkflowSignalError::Serialization(_) => None,
1252 }
1253 }
1254
1255 pub fn reason(&self) -> Option<&IncomingError> {
1257 match self {
1258 WorkflowSignalError::NotFound(err) | WorkflowSignalError::Failed(err) => {
1259 Some(err.error())
1260 }
1261 WorkflowSignalError::Serialization(_) => None,
1262 }
1263 }
1264
1265 pub fn as_cancelled(&self) -> Option<&CancelledError> {
1268 self.reason()?.as_cancelled()
1269 }
1270}
1271
1272#[derive(Debug)]
1274pub struct WorkflowSignalFailureError {
1275 failure: Failure,
1276 error: Box<IncomingError>,
1277}
1278
1279impl WorkflowSignalFailureError {
1280 pub(crate) fn new(failure: Failure, error: IncomingError) -> Self {
1282 Self {
1283 failure,
1284 error: Box::new(error),
1285 }
1286 }
1287
1288 pub fn failure(&self) -> &Failure {
1290 &self.failure
1291 }
1292
1293 pub fn cause(&self) -> Option<&IncomingError> {
1295 self.error.cause()
1296 }
1297
1298 pub fn error(&self) -> &IncomingError {
1300 &self.error
1301 }
1302}
1303
1304impl std::fmt::Display for WorkflowSignalFailureError {
1305 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1306 self.failure.fmt(f)
1307 }
1308}
1309
1310impl std::error::Error for WorkflowSignalFailureError {
1311 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
1312 self.cause()
1313 .map(|cause| cause as &(dyn std::error::Error + 'static))
1314 }
1315}
1316
1317#[derive(Debug)]
1319pub struct WorkflowCancelFailureError {
1320 failure: Failure,
1321 error: Box<IncomingError>,
1322}
1323
1324impl WorkflowCancelFailureError {
1325 pub(crate) fn new(failure: Failure, error: IncomingError) -> Self {
1327 Self {
1328 failure,
1329 error: Box::new(error),
1330 }
1331 }
1332
1333 pub fn failure(&self) -> &Failure {
1335 &self.failure
1336 }
1337
1338 pub fn cause(&self) -> Option<&IncomingError> {
1340 self.error.cause()
1341 }
1342
1343 pub fn error(&self) -> &IncomingError {
1345 &self.error
1346 }
1347}
1348
1349impl std::fmt::Display for WorkflowCancelFailureError {
1350 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1351 self.failure.fmt(f)
1352 }
1353}
1354
1355impl std::error::Error for WorkflowCancelFailureError {
1356 fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
1357 self.cause()
1358 .map(|cause| cause as &(dyn std::error::Error + 'static))
1359 }
1360}
1361
1362#[cfg(test)]
1363mod tests {
1364 use super::*;
1365 use crate::{
1366 data_converters::{
1367 DefaultFailureConverter, FailureConverter, GenericPayloadConverter, PayloadConverter,
1368 SerializationContext, SerializationContextData, WorkflowSerializationContext,
1369 },
1370 protos::temporal::api::{
1371 common::v1::Payload,
1372 failure::v1::{ActivityFailureInfo, TimeoutFailureInfo, failure::FailureInfo},
1373 },
1374 };
1375
1376 struct AlwaysFailsSerialize;
1377
1378 impl serde::Serialize for AlwaysFailsSerialize {
1379 fn serialize<S: serde::Serializer>(&self, _serializer: S) -> Result<S::Ok, S::Error> {
1380 Err(serde::ser::Error::custom("serialize boom"))
1381 }
1382 }
1383
1384 #[test]
1385 fn decoded_failures_hide_unknown_state_values() {
1386 let failure = Failure {
1387 cause: Some(Box::new(Failure {
1388 failure_info: Some(FailureInfo::TimeoutFailureInfo(TimeoutFailureInfo {
1389 timeout_type: 654_321,
1390 ..Default::default()
1391 })),
1392 ..Default::default()
1393 })),
1394 failure_info: Some(FailureInfo::ActivityFailureInfo(ActivityFailureInfo {
1395 retry_state: 123_456,
1396 ..Default::default()
1397 })),
1398 ..Default::default()
1399 };
1400
1401 let decoded = DefaultFailureConverter::default()
1402 .to_error(
1403 failure,
1404 &PayloadConverter::default(),
1405 &SerializationContextData::Workflow(WorkflowSerializationContext::new()),
1406 )
1407 .unwrap();
1408 let IncomingError::Activity(activity) = decoded else {
1409 panic!("expected activity failure");
1410 };
1411 assert_eq!(activity.retry_state(), RetryState::Unknown);
1412 let Some(FailureInfo::ActivityFailureInfo(raw_activity)) =
1413 activity.failure().failure_info.as_ref()
1414 else {
1415 panic!("expected raw activity failure info");
1416 };
1417 assert_eq!(raw_activity.retry_state, 123_456);
1418 let Some(IncomingError::Timeout(timeout)) = activity.cause() else {
1419 panic!("expected timeout cause");
1420 };
1421 assert_eq!(timeout.timeout_type(), TimeoutType::Unknown);
1422 let Some(FailureInfo::TimeoutFailureInfo(raw_timeout)) =
1423 timeout.failure().failure_info.as_ref()
1424 else {
1425 panic!("expected raw timeout failure info");
1426 };
1427 assert_eq!(raw_timeout.timeout_type, 654_321);
1428 }
1429
1430 #[test]
1431 fn constructors_set_retryability_defaults() {
1432 assert!(!ApplicationFailure::new(anyhow::anyhow!("retryable")).is_non_retryable());
1433 assert!(
1434 ApplicationFailure::non_retryable(anyhow::anyhow!("non-retryable")).is_non_retryable()
1435 );
1436 }
1437
1438 #[test]
1439 fn conversion_preserves_application_metadata() {
1440 let payloads = Payloads {
1441 payloads: vec![Payload {
1442 data: b"details".to_vec(),
1443 ..Default::default()
1444 }],
1445 };
1446 let failure = DefaultFailureConverter::default().to_failure(
1447 OutgoingError::Workflow(OutgoingWorkflowError::Application(Box::new(
1448 ApplicationFailure::builder(anyhow::anyhow!("oops"))
1449 .type_name("MyType".to_owned())
1450 .non_retryable(true)
1451 .next_retry_delay(Duration::from_secs(3))
1452 .category(ApplicationErrorCategory::Benign)
1453 .details(RawValue::new(payloads.payloads.clone()))
1454 .build(),
1455 ))),
1456 &PayloadConverter::default(),
1457 &SerializationContextData::None,
1458 );
1459 let Some(FailureInfo::ApplicationFailureInfo(info)) = failure.failure_info else {
1460 panic!("expected application failure info");
1461 };
1462 assert_eq!(failure.message, "oops");
1463 assert_eq!(info.r#type, "MyType");
1464 assert!(info.non_retryable);
1465 assert_eq!(info.details, Some(payloads));
1466 assert_eq!(info.category(), ProtoApplicationErrorCategory::Benign);
1467 assert_eq!(info.next_retry_delay.unwrap().seconds, 3);
1468 }
1469
1470 #[test]
1471 fn builder_accepts_raw_payload_details() {
1472 let payload = Payload {
1473 data: b"details".to_vec(),
1474 ..Default::default()
1475 };
1476 let failure = DefaultFailureConverter::default().to_failure(
1477 OutgoingError::Workflow(OutgoingWorkflowError::Application(Box::new(
1478 ApplicationFailure::builder(anyhow::anyhow!("oops"))
1479 .details(RawValue::new(vec![payload.clone()]))
1480 .build(),
1481 ))),
1482 &PayloadConverter::default(),
1483 &SerializationContextData::None,
1484 );
1485
1486 let Some(FailureInfo::ApplicationFailureInfo(info)) = failure.failure_info else {
1487 panic!("expected application failure info");
1488 };
1489 assert_eq!(info.details.unwrap().payloads, vec![payload]);
1490 }
1491
1492 #[test]
1493 fn builder_accepts_serializable_details() {
1494 let failure = DefaultFailureConverter::default().to_failure(
1495 OutgoingError::Workflow(OutgoingWorkflowError::Application(Box::new(
1496 ApplicationFailure::builder(anyhow::anyhow!("oops"))
1497 .details("details".to_string())
1498 .build(),
1499 ))),
1500 &PayloadConverter::default(),
1501 &SerializationContextData::None,
1502 );
1503
1504 let Some(FailureInfo::ApplicationFailureInfo(info)) = failure.failure_info else {
1505 panic!("expected application failure info");
1506 };
1507 let payloads = info.details.expect("expected details").payloads;
1508 let converter = PayloadConverter::default();
1509 let details: String = converter
1510 .from_payloads(
1511 &SerializationContext::new(&SerializationContextData::None, &converter),
1512 payloads,
1513 )
1514 .unwrap();
1515 assert_eq!(details, "details");
1516 }
1517
1518 #[test]
1519 fn application_failure_encoding_surfaces_detail_encoding_errors() {
1520 let failure = DefaultFailureConverter::default().to_failure(
1521 OutgoingError::Workflow(OutgoingWorkflowError::Application(Box::new(
1522 ApplicationFailure::builder(anyhow::anyhow!("oops"))
1523 .details(AlwaysFailsSerialize)
1524 .build(),
1525 ))),
1526 &PayloadConverter::default(),
1527 &SerializationContextData::None,
1528 );
1529
1530 assert_eq!(
1531 failure.message,
1532 "Failed converting error to failure: Encoding error: serialize boom, original error message: oops"
1533 );
1534 assert!(matches!(
1535 failure.failure_info,
1536 Some(FailureInfo::ApplicationFailureInfo(_))
1537 ));
1538 }
1539
1540 #[test]
1541 fn anyhow_workflow_errors_default_to_application_outgoing_errors() {
1542 let outgoing: OutgoingWorkflowError = anyhow::anyhow!("workflow boom").into();
1543
1544 let OutgoingWorkflowError::Application(app) = outgoing else {
1545 panic!("plain workflow errors should default to application failures");
1546 };
1547 assert_eq!(app.to_string(), "workflow boom");
1548 }
1549
1550 #[test]
1551 fn payload_conversion_errors_use_dedicated_outgoing_variant() {
1552 let outgoing: OutgoingWorkflowError =
1553 PayloadConversionError::EncodingError(anyhow::anyhow!("encode boom").into()).into();
1554
1555 let OutgoingWorkflowError::PayloadConversion(err) = outgoing else {
1556 panic!("expected a payload conversion failure");
1557 };
1558 assert_eq!(err.to_string(), "Encoding error: encode boom");
1559 }
1560}