( effect: Effect.Effect<A, E, R>, isSuspend: (value: A) => boolean )
| 717 | * @since 4.0.0 |
| 718 | */ |
| 719 | export const wrapActivityResult = <A, E, R>( |
| 720 | effect: Effect.Effect<A, E, R>, |
| 721 | isSuspend: (value: A) => boolean |
| 722 | ): Effect.Effect<A, E, R | WorkflowInstance> => |
| 723 | Effect.contextWith((context: Context.Context<WorkflowInstance>) => { |
| 724 | const instance = Context.get(context, InstanceTag) |
| 725 | const state = instance.activityState |
| 726 | if (state.count === 0) state.latch.closeUnsafe() |
| 727 | state.count++ |
| 728 | return Effect.onExit(effect, (exit) => { |
| 729 | state.count-- |
| 730 | const isSuspended = Exit.isSuccess(exit) && isSuspend(exit.value) |
| 731 | if ( |
| 732 | Exit.isSuccess(exit) && |
| 733 | isResult(exit.value) && |
| 734 | exit.value._tag === "Suspended" && |
| 735 | exit.value.cause |
| 736 | ) { |
| 737 | instance.cause = instance.cause |
| 738 | ? Cause.combine(instance.cause, exit.value.cause) |
| 739 | : exit.value.cause |
| 740 | } |
| 741 | return state.count === 0 |
| 742 | ? state.latch.open |
| 743 | : isSuspended |
| 744 | ? waitForZero(instance) |
| 745 | : Effect.void |
| 746 | }) |
| 747 | }) |
| 748 | |
| 749 | const waitForZero = Effect.fnUntraced(function*(instance: WorkflowInstance["Service"]) { |
| 750 | const state = instance.activityState |
nothing calls this directly
no test coverage detected