MCPcopy Create free account
hub / github.com/effect-app/libs / collectUnprocessedById

Function collectUnprocessedById

packages/infra/src/ClusterCosmos.ts:743–763  ·  view source on GitHub ↗
(
  docs: ReadonlyArray<MessageDoc>,
  queryReplies: (
    query: string,
    parameters: ReadonlyArray<CosmosParameter>
  ) => Effect.Effect<Array<ReplyDoc>, E>
)

Source from the content-addressed store, hash-verified

741 })
742
743const collectUnprocessedById = <E>(
744 docs: ReadonlyArray<MessageDoc>,
745 queryReplies: (
746 query: string,
747 parameters: ReadonlyArray<CosmosParameter>
748 ) => Effect.Effect<Array<ReplyDoc>, E>
749) =>
750 Effect.gen(function*() {
751 const messages: Array<{
752 readonly envelope: Envelope.Encoded
753 readonly lastSentReply: Option.Option<Reply.Encoded>
754 }> = []
755 const activeRequestIds = yield* activeReplyRequestIds(docs, queryReplies)
756 const lastReplies = yield* lastRepliesById(docs, activeRequestIds, queryReplies)
757 for (const doc of docs) {
758 if (activeRequestIds.has(doc.requestId)) continue
759 const sentReply = Option.fromNullishOr(doc.lastReplyId === null ? undefined : lastReplies.get(doc.lastReplyId))
760 messages.push(envelopeFromDoc(doc, sentReply))
761 }
762 return messages
763 })
764
765const activeReplyRequestIds = <E>(
766 docs: ReadonlyArray<MessageDoc>,

Callers 1

ClusterCosmos.tsFile · 0.85

Calls 6

activeReplyRequestIdsFunction · 0.85
lastRepliesByIdFunction · 0.85
envelopeFromDocFunction · 0.85
pushMethod · 0.80
getMethod · 0.65
hasMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…