(
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>
)
| 714 | }) |
| 715 | |
| 716 | const 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 | |
| 743 | const collectUnprocessedById = <E>( |
| 744 | docs: ReadonlyArray<MessageDoc>, |
no test coverage detected
searching dependent graphs…