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

Function decodeReplies

packages/effect/src/unstable/cluster/MessageStorage.ts:741–784  ·  view source on GitHub ↗
(
    messages: Map<string, Message.OutgoingRequest<any>>,
    encodedReplies: Array<Reply.Encoded>
  )

Source from the content-addressed store, hash-verified

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})

Callers 1

MessageStorage.tsFile · 0.85

Calls 6

pushMethod · 0.80
getMethod · 0.65
pipeMethod · 0.65
addMethod · 0.65
hasMethod · 0.45
succeedMethod · 0.45

Tested by

no test coverage detected