(
storage: MessageStorage["Service"],
envelopes: Array<{
readonly envelope: Envelope.Encoded
readonly lastSentReply: Option.Option<Reply.Encoded>
}>
)
| 685 | }) |
| 686 | |
| 687 | const decodeMessages = ( |
| 688 | storage: MessageStorage["Service"], |
| 689 | envelopes: Array<{ |
| 690 | readonly envelope: Envelope.Encoded |
| 691 | readonly lastSentReply: Option.Option<Reply.Encoded> |
| 692 | }> |
| 693 | ) => { |
| 694 | const messages: Array<Message.Incoming<any>> = [] |
| 695 | let index = 0 |
| 696 | |
| 697 | // if we have a malformed message, we should not return it and update |
| 698 | // the storage with a defect |
| 699 | const decodeMessage = Effect.catch( |
| 700 | Effect.suspend(() => { |
| 701 | const envelope = envelopes[index] |
| 702 | if (!envelope) return Effect.undefined |
| 703 | return decodeEnvelopeWithReply(envelope) |
| 704 | }), |
| 705 | (error) => { |
| 706 | const envelope = envelopes[index] |
| 707 | return storage.saveReply(Reply.ReplyWithContext.fromDefect({ |
| 708 | id: snowflakeGen.nextUnsafe(), |
| 709 | requestId: Snowflake.Snowflake(envelope.envelope.requestId), |
| 710 | defect: error.toString() |
| 711 | })).pipe( |
| 712 | Effect.forkDetach, |
| 713 | Effect.asVoid |
| 714 | ) |
| 715 | } |
| 716 | ) |
| 717 | return Effect.as( |
| 718 | Effect.whileLoop({ |
| 719 | while: () => index < envelopes.length, |
| 720 | body: () => decodeMessage, |
| 721 | step: (message) => { |
| 722 | const envelope = envelopes[index++] |
| 723 | if (!message) return |
| 724 | messages.push( |
| 725 | message.envelope._tag === "Request" |
| 726 | ? new Message.IncomingRequest({ |
| 727 | envelope: message.envelope, |
| 728 | lastSentReply: envelope.lastSentReply, |
| 729 | respond: storage.saveReply |
| 730 | }) |
| 731 | : new Message.IncomingEnvelope({ |
| 732 | envelope: message.envelope |
| 733 | }) |
| 734 | ) |
| 735 | } |
| 736 | }), |
| 737 | messages |
| 738 | ) |
| 739 | } |
| 740 | |
| 741 | const decodeReplies = ( |
| 742 | messages: Map<string, Message.OutgoingRequest<any>>, |
no test coverage detected