| 739 | } |
| 740 | |
| 741 | const decodeReplies = ( |
| 742 | messages: Map<string, Message.OutgoingRequest<any>>, |
| 743 | encodedReplies: Array<Reply.Encoded> |
| 744 | ) => { |
| 745 | const replies: Array<Reply.Reply<any>> = [] |
| 746 | const ignoredRequests = new Set<string>() |
| 747 | let index = 0 |
| 748 | |
| 749 | const decodeReply: Effect.Effect<void | Reply.Reply<any>> = Effect.catch( |
| 750 | Effect.suspend(() => { |
| 751 | const reply = encodedReplies[index] |
| 752 | if (ignoredRequests.has(reply.requestId)) return Effect.void |
| 753 | const message = messages.get(reply.requestId) |
| 754 | if (!message) return Effect.void |
| 755 | const schema = Reply.Reply(message.rpc) |
| 756 | return Schema.decodeEffect(schema)(reply).pipe( |
| 757 | Effect.provideContext(message.context) |
| 758 | ) as Effect.Effect<Reply.Reply<any>, Schema.SchemaError> |
| 759 | }), |
| 760 | (error) => { |
| 761 | const reply = encodedReplies[index] |
| 762 | ignoredRequests.add(reply.requestId) |
| 763 | return Effect.succeed( |
| 764 | new Reply.WithExit({ |
| 765 | id: snowflakeGen.nextUnsafe(), |
| 766 | requestId: Snowflake.Snowflake(reply.requestId), |
| 767 | exit: Exit.die(error) |
| 768 | }) |
| 769 | ) |
| 770 | } |
| 771 | ) |
| 772 | |
| 773 | return Effect.as( |
| 774 | Effect.whileLoop({ |
| 775 | while: () => index < encodedReplies.length, |
| 776 | body: () => decodeReply, |
| 777 | step: (reply) => { |
| 778 | index++ |
| 779 | if (reply) replies.push(reply) |
| 780 | } |
| 781 | }), |
| 782 | replies |
| 783 | ) |
| 784 | } |
| 785 | |
| 786 | return storage |
| 787 | }) |