( state: ActiveReplay, handler: (value: A) => Effect.Effect<unknown, E, R> | void, decode: (event: WebSocketEvent) => A, onOpen: Effect.Effect<void> | undefined, )
| 72 | }) |
| 73 | |
| 74 | const runReplay = <A, E, R>( |
| 75 | state: ActiveReplay, |
| 76 | handler: (value: A) => Effect.Effect<unknown, E, R> | void, |
| 77 | decode: (event: WebSocketEvent) => A, |
| 78 | onOpen: Effect.Effect<void> | undefined, |
| 79 | ) => |
| 80 | Effect.scoped( |
| 81 | Effect.gen(function* () { |
| 82 | const handlers = yield* FiberSet.make<unknown, E>() |
| 83 | const run = yield* FiberSet.runtime(handlers)<R>() |
| 84 | if (onOpen) yield* onOpen |
| 85 | |
| 86 | const drive = Effect.gen(function* () { |
| 87 | while (true) { |
| 88 | const current = yield* Ref.get(state.progress) |
| 89 | const event = state.interaction.events[current.position] |
| 90 | if (!event) return |
| 91 | if (yield* Ref.get(state.closed)) |
| 92 | return yield* Effect.die( |
| 93 | new Error( |
| 94 | `WebSocket closed with unconsumed events: used ${current.position} of ${state.interaction.events.length}`, |
| 95 | ), |
| 96 | ) |
| 97 | if (event.direction === "server") { |
| 98 | yield* Ref.set(state.progress, { |
| 99 | position: current.position + 1, |
| 100 | changed: yield* Deferred.make<void>(), |
| 101 | }) |
| 102 | run(runHandler(handler, decode(event))) |
| 103 | continue |
| 104 | } |
| 105 | yield* Deferred.await(current.changed) |
| 106 | } |
| 107 | }) |
| 108 | |
| 109 | yield* drive.pipe(Effect.raceFirst(FiberSet.join(handlers))) |
| 110 | yield* FiberSet.awaitEmpty(handlers).pipe(Effect.raceFirst(FiberSet.join(handlers))) |
| 111 | }), |
| 112 | ) |
| 113 | |
| 114 | const openSnapshot = (request: WebSocketRequest, redactor: Redactor) => { |
| 115 | const snapshot = redactor.request({ method: "GET", url: request.url, headers: request.headers ?? {}, body: "" }) |
no test coverage detected