| 172 | } |
| 173 | |
| 174 | function wrapCommandHandle<E extends ExecEncoding>( |
| 175 | handle: ExecHandle<E>, |
| 176 | backend: string, |
| 177 | lazyResultHandle?: () => Promise<ExecHandle<E>>, |
| 178 | ): WorkspaceRuntimeExecHandle<E> { |
| 179 | let claimed: "result" | "stream" | undefined; |
| 180 | let reader: ReadableStreamDefaultReader<WorkspaceRuntimeEvent<E>> | undefined; |
| 181 | let resultPromise: Promise<WorkspaceRuntimeResult<E>> | undefined; |
| 182 | const stream = new ReadableStream<WorkspaceRuntimeEvent<E>>( |
| 183 | { |
| 184 | async pull(controller) { |
| 185 | if (claimed === "result") { |
| 186 | controller.error(new Error("runtime handle already consumed by result()")); |
| 187 | return; |
| 188 | } |
| 189 | claimed = "stream"; |
| 190 | reader ??= handle.getReader() as ReadableStreamDefaultReader<WorkspaceRuntimeEvent<E>>; |
| 191 | try { |
| 192 | const next = await reader.read(); |
| 193 | if (next.done) { |
| 194 | reader.releaseLock(); |
| 195 | reader = undefined; |
| 196 | controller.close(); |
| 197 | } else controller.enqueue(next.value); |
| 198 | } catch (error) { |
| 199 | reader?.releaseLock(); |
| 200 | reader = undefined; |
| 201 | controller.error(error); |
| 202 | } |
| 203 | }, |
| 204 | async cancel(reason) { |
| 205 | if (reader) { |
| 206 | try { |
| 207 | await reader.cancel(reason); |
| 208 | } finally { |
| 209 | reader.releaseLock(); |
| 210 | reader = undefined; |
| 211 | } |
| 212 | } else await handle.cancel(reason); |
| 213 | }, |
| 214 | }, |
| 215 | { highWaterMark: 0 }, |
| 216 | ) as WorkspaceRuntimeExecHandle<E>; |
| 217 | Object.defineProperties(stream, { |
| 218 | id: { value: handle.id, enumerable: false }, |
| 219 | backend: { value: backend, enumerable: false }, |
| 220 | result: { |
| 221 | value: (): Promise<WorkspaceRuntimeResult<E>> => { |
| 222 | if (claimed === "stream") { |
| 223 | throw new Error("runtime handle already streaming: result() and streaming are exclusive"); |
| 224 | } |
| 225 | claimed = "result"; |
| 226 | resultPromise ??= (async () => { |
| 227 | if (lazyResultHandle) await handle.cancel("result() requested a full replay"); |
| 228 | const result = lazyResultHandle |
| 229 | ? await (await lazyResultHandle()).result() |
| 230 | : await handle.result(); |
| 231 | return { |