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

Function collectUnprocessed

packages/infra/src/ClusterCosmos.ts:716–741  ·  view source on GitHub ↗
(
  docs: ReadonlyArray<MessageDoc>,
  now: number,
  claimMessageRead: (doc: MessageDoc, now: number) => Effect.Effect<boolean, E>,
  queryReplies: (
    query: string,
    parameters: ReadonlyArray<CosmosParameter>
  ) => Effect.Effect<Array<ReplyDoc>, E>
)

Source from the content-addressed store, hash-verified

714})
715
716const collectUnprocessed = <E>(
717 docs: ReadonlyArray<MessageDoc>,
718 now: number,
719 claimMessageRead: (doc: MessageDoc, now: number) => Effect.Effect<boolean, E>,
720 queryReplies: (
721 query: string,
722 parameters: ReadonlyArray<CosmosParameter>
723 ) => Effect.Effect<Array<ReplyDoc>, E>
724) =>
725 Effect.gen(function*() {
726 const messages: Array<{
727 readonly envelope: Envelope.Encoded
728 readonly lastSentReply: Option.Option<Reply.Encoded>
729 }> = []
730 const activeRequestIds = yield* activeReplyRequestIds(docs, queryReplies)
731 const lastReplies = yield* lastRepliesById(docs, activeRequestIds, queryReplies)
732 for (const doc of docs) {
733 if (activeRequestIds.has(doc.requestId)) continue
734 const sentReply = Option.fromNullishOr(doc.lastReplyId === null ? undefined : lastReplies.get(doc.lastReplyId))
735 const claimed = yield* claimMessageRead(doc, now)
736 if (claimed) {
737 messages.push(envelopeFromDoc({ ...doc, lastRead: now }, sentReply))
738 }
739 }
740 return messages
741 })
742
743const collectUnprocessedById = <E>(
744 docs: ReadonlyArray<MessageDoc>,

Callers 1

ClusterCosmos.tsFile · 0.85

Calls 7

activeReplyRequestIdsFunction · 0.85
lastRepliesByIdFunction · 0.85
claimMessageReadFunction · 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…