(stream: ReadableStream<T>)
| 37 | } |
| 38 | |
| 39 | async function drainStream<T>(stream: ReadableStream<T>): Promise<T[]> { |
| 40 | const out: T[] = []; |
| 41 | const reader = stream.getReader(); |
| 42 | try { |
| 43 | while (true) { |
| 44 | const { value, done } = await reader.read(); |
| 45 | if (done) break; |
| 46 | out.push(value); |
| 47 | } |
| 48 | } finally { |
| 49 | reader.releaseLock(); |
| 50 | } |
| 51 | return out; |
| 52 | } |
| 53 | |
| 54 | // Wrap a SyncRPC so the fetchChanges result carries a tracked |
| 55 | // [Symbol.dispose]. pullOnce owns that envelope and must dispose it on |
no test coverage detected