(
resolver: RequestResolver<A>,
request: A,
resume: (exit: Exit<any, any>) => void,
fiber: {
readonly context: Context.Context<never>
readonly currentScheduler: Scheduler
readonly id?: number
}
)
| 82 | const pendingBatches = new WeakMap<RequestResolver<any>, Map<unknown, Batch>>() |
| 83 | |
| 84 | const addEntry = <A extends Request.Any>( |
| 85 | resolver: RequestResolver<A>, |
| 86 | request: A, |
| 87 | resume: (exit: Exit<any, any>) => void, |
| 88 | fiber: { |
| 89 | readonly context: Context.Context<never> |
| 90 | readonly currentScheduler: Scheduler |
| 91 | readonly id?: number |
| 92 | } |
| 93 | ) => { |
| 94 | let batchMap = pendingBatches.get(resolver) |
| 95 | if (!batchMap) { |
| 96 | batchMap = new Map<object, Batch>() |
| 97 | pendingBatches.set(resolver, batchMap) |
| 98 | } |
| 99 | let batch: Batch | undefined |
| 100 | let completed = false |
| 101 | const entry = makeEntry({ |
| 102 | request, |
| 103 | context: fiber.context as any, |
| 104 | uninterruptible: false, |
| 105 | completeUnsafe(effect) { |
| 106 | if (completed) return |
| 107 | completed = true |
| 108 | resume(effect) |
| 109 | batch?.entrySet.delete(entry) |
| 110 | } |
| 111 | }) |
| 112 | if (resolver.preCheck !== undefined && !resolver.preCheck(entry)) { |
| 113 | return entry |
| 114 | } |
| 115 | const key = resolver.batchKey(entry) |
| 116 | batch = batchMap.get(key) |
| 117 | if (!batch) { |
| 118 | if (batchPool.length > 0) { |
| 119 | batch = batchPool.pop()! |
| 120 | batch.key = key |
| 121 | batch.resolver = resolver |
| 122 | batch.map = batchMap |
| 123 | } else { |
| 124 | const newBatch: Batch = { |
| 125 | key, |
| 126 | resolver, |
| 127 | map: batchMap, |
| 128 | entrySet: new Set(), |
| 129 | entries: new Set(), |
| 130 | delayEffect: effect.flatMap( |
| 131 | effect.suspend(() => newBatch.resolver.delay), |
| 132 | (_) => runBatch(newBatch) |
| 133 | ) as Effect<void>, |
| 134 | run: effect.onExit( |
| 135 | effect.suspend(() => |
| 136 | newBatch.resolver.runAll(Array.from(newBatch.entries) as NonEmptyArray<Request.Entry<any>>, newBatch.key) |
| 137 | ), |
| 138 | (exit) => { |
| 139 | for (const entry of newBatch.entrySet) { |
| 140 | entry.completeUnsafe( |
| 141 | exit._tag === "Success" |
no test coverage detected