MCPcopy Create free account
hub / github.com/Effect-TS/effect / decodeMessages

Function decodeMessages

packages/effect/src/unstable/cluster/MessageStorage.ts:687–739  ·  view source on GitHub ↗
(
    storage: MessageStorage["Service"],
    envelopes: Array<{
      readonly envelope: Envelope.Encoded
      readonly lastSentReply: Option.Option<Reply.Encoded>
    }>
  )

Source from the content-addressed store, hash-verified

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>>,

Callers 2

unprocessedMessagesFunction · 0.85
unprocessedMessagesByIdFunction · 0.85

Calls 4

fromDefectMethod · 0.80
pushMethod · 0.80
pipeMethod · 0.65
toStringMethod · 0.65

Tested by

no test coverage detected