()
| 37 | } |
| 38 | |
| 39 | function createEventStream() { |
| 40 | const queue: GlobalEventEnvelope[] = [] |
| 41 | const waiters: Array<(value: GlobalEventEnvelope | undefined) => void> = [] |
| 42 | const state = { closed: false } |
| 43 | |
| 44 | const push = (event: GlobalEventEnvelope) => { |
| 45 | const waiter = waiters.shift() |
| 46 | if (waiter) { |
| 47 | waiter(event) |
| 48 | return |
| 49 | } |
| 50 | queue.push(event) |
| 51 | } |
| 52 | |
| 53 | const close = () => { |
| 54 | state.closed = true |
| 55 | for (const waiter of waiters.splice(0)) { |
| 56 | waiter(undefined) |
| 57 | } |
| 58 | } |
| 59 | |
| 60 | const stream = async function* (signal?: AbortSignal) { |
| 61 | while (true) { |
| 62 | if (signal?.aborted) return |
| 63 | const next = queue.shift() |
| 64 | if (next) { |
| 65 | yield next |
| 66 | continue |
| 67 | } |
| 68 | if (state.closed) return |
| 69 | const value = await new Promise<GlobalEventEnvelope | undefined>((resolve) => { |
| 70 | waiters.push(resolve) |
| 71 | signal?.addEventListener("abort", () => resolve(undefined), { once: true }) |
| 72 | }) |
| 73 | if (!value) return |
| 74 | yield value |
| 75 | } |
| 76 | } |
| 77 | |
| 78 | return { push, close, stream } |
| 79 | } |
| 80 | |
| 81 | function createHarness(messages: Record<string, SessionMessageResponse> = {}) { |
| 82 | const updates: SessionUpdateParams[] = [] |
no test coverage detected