(options?: DecodeOptions)
| 100 | * @since 4.0.0 |
| 101 | */ |
| 102 | export const decode = <IE, Done>(options?: DecodeOptions): Channel.Channel< |
| 103 | NonEmptyReadonlyArray<Event>, |
| 104 | IE | Retry | SseError, |
| 105 | Done, |
| 106 | NonEmptyReadonlyArray<string>, |
| 107 | IE, |
| 108 | Done |
| 109 | > => |
| 110 | Channel.fromTransform((upstream, _scope) => |
| 111 | Effect.sync(() => { |
| 112 | let buffer: Array<Event> = [] |
| 113 | let retry: Retry | undefined |
| 114 | const parser = makeParser((event) => { |
| 115 | if (event._tag === "Retry") { |
| 116 | retry = event |
| 117 | } else { |
| 118 | buffer.push(event) |
| 119 | } |
| 120 | }, options) |
| 121 | |
| 122 | const pump = Effect.flatMap(upstream, (arr) => { |
| 123 | for (let i = 0; i < arr.length; i++) { |
| 124 | const error = parser.feed(arr[i]) |
| 125 | if (error !== undefined) { |
| 126 | return Effect.fail(error) |
| 127 | } |
| 128 | } |
| 129 | return Effect.void |
| 130 | }) |
| 131 | |
| 132 | return Effect.suspend(function loop(): Pull.Pull<NonEmptyReadonlyArray<Event>, IE | Retry | SseError, Done> { |
| 133 | if (Arr.isArrayNonEmpty(buffer)) { |
| 134 | const out = buffer |
| 135 | buffer = [] |
| 136 | return Effect.succeed(out) |
| 137 | } else if (retry) { |
| 138 | return Effect.fail(retry) |
| 139 | } |
| 140 | return Effect.flatMap(pump, loop) |
| 141 | }) |
| 142 | }) |
| 143 | ) |
| 144 | |
| 145 | /** |
| 146 | * A constraint for schemas that can decode SSE events. |
no test coverage detected